use std::collections::{BTreeSet, HashMap, HashSet};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use fsqlite_types::{CommitSeq, IntentOp, ObjectId, PageData, PageNumber, Snapshot, TxnToken};
use parking_lot::RwLock;
use tracing::{debug, info, warn};
use crate::core_types::TransactionMode;
use crate::witness_objects::AbortPolicy;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CoordinatorMode {
Native,
Compatibility,
}
#[derive(Debug)]
pub struct NativePublishRequest {
pub txn: TxnToken,
pub begin_seq: CommitSeq,
pub capsule_object_id: ObjectId,
pub capsule_digest: [u8; 32],
pub write_set_summary: BTreeSet<u32>,
pub read_witnesses: Vec<ObjectId>,
pub write_witnesses: Vec<ObjectId>,
pub edge_ids: Vec<ObjectId>,
pub merge_witnesses: Vec<ObjectId>,
pub abort_policy: AbortPolicy,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum NativePublishResponse {
Ok {
commit_seq: CommitSeq,
marker_object_id: ObjectId,
},
Conflict {
conflicting_pages: Vec<PageNumber>,
conflicting_commit_seq: CommitSeq,
},
Aborted {
code: u32,
},
IoError {
message: String,
},
}
#[derive(Debug)]
pub struct CompatCommitRequest {
pub txn: TxnToken,
pub mode: TransactionMode,
pub write_set: CommitWriteSet,
pub intent_log: Vec<IntentOp>,
pub page_locks: HashSet<PageNumber>,
pub snapshot: Snapshot,
pub has_in_rw: bool,
pub has_out_rw: bool,
pub wal_fec_r: u8,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CompatCommitResponse {
Ok {
wal_offset: u64,
commit_seq: CommitSeq,
},
Conflict {
conflicting_pages: Vec<PageNumber>,
conflicting_commit_seq: CommitSeq,
},
IoError {
message: String,
},
}
#[derive(Debug)]
pub enum CommitWriteSet {
Inline(HashMap<PageNumber, PageData>),
Spilled(SpilledWriteSet),
}
impl CommitWriteSet {
#[must_use]
pub fn page_count(&self) -> usize {
match self {
Self::Inline(pages) => pages.len(),
Self::Spilled(spilled) => spilled.pages.len(),
}
}
#[must_use]
pub fn page_numbers(&self) -> Vec<PageNumber> {
match self {
Self::Inline(pages) => pages.keys().copied().collect(),
Self::Spilled(spilled) => spilled.pages.keys().copied().collect(),
}
}
#[must_use]
pub const fn is_spilled(&self) -> bool {
matches!(self, Self::Spilled(_))
}
}
#[derive(Debug)]
pub enum SpillHandle {
Path(PathBuf),
#[cfg(target_family = "unix")]
Fd(std::os::unix::io::OwnedFd),
}
#[derive(Debug)]
pub struct SpilledWriteSet {
pub spill: SpillHandle,
pub pages: HashMap<PageNumber, SpillLoc>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SpillLoc {
pub offset: u64,
pub len: u32,
pub xxh3_64: u64,
}
pub const DEFAULT_SPILL_THRESHOLD: usize = 32 * 1024 * 1024;
pub const DEFAULT_MAX_BATCH_SIZE: usize = 16;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct CoordinatorLease {
pub holder_pid: u64,
pub acquired_at: u64,
pub expires_at: u64,
}
pub struct WriteCoordinator {
mode: CoordinatorMode,
next_commit_seq: AtomicU64,
commit_page_index: RwLock<HashMap<u32, CommitSeq>>,
wal_offset: AtomicU64,
lease: RwLock<Option<CoordinatorLease>>,
committed_seqs: RwLock<Vec<CommitSeq>>,
}
impl WriteCoordinator {
#[must_use]
pub fn new(mode: CoordinatorMode) -> Self {
Self {
mode,
next_commit_seq: AtomicU64::new(1),
commit_page_index: RwLock::new(HashMap::new()),
wal_offset: AtomicU64::new(0),
lease: RwLock::new(None),
committed_seqs: RwLock::new(Vec::new()),
}
}
#[must_use]
pub fn mode(&self) -> CoordinatorMode {
self.mode
}
pub fn restore_state(
&self,
next_seq: CommitSeq,
recent_commits: HashMap<u32, CommitSeq>,
wal_offset: u64,
) {
self.next_commit_seq.store(next_seq.get(), Ordering::SeqCst);
self.wal_offset.store(wal_offset, Ordering::SeqCst);
let mut index = self.commit_page_index.write();
*index = recent_commits;
info!(
bead_id = "bd-389e",
next_seq = next_seq.get(),
wal_offset,
restored_pages = index.len(),
"coordinator state restored from persistence"
);
}
pub fn acquire_lease(&self, pid: u64, timestamp: u64) -> bool {
let mut lease = self.lease.write();
if let Some(existing) = &*lease {
if existing.expires_at > 0 && existing.expires_at <= timestamp {
info!(
bead_id = "bd-389e",
old_pid = existing.holder_pid,
new_pid = pid,
"coordinator lease expired, allowing takeover"
);
} else {
debug!(
bead_id = "bd-389e",
holder = existing.holder_pid,
"coordinator lease already held"
);
return false;
}
}
*lease = Some(CoordinatorLease {
holder_pid: pid,
acquired_at: timestamp,
expires_at: 0, });
drop(lease);
info!(bead_id = "bd-389e", pid, "coordinator lease acquired");
true
}
pub fn release_lease(&self, pid: u64) -> bool {
let mut lease = self.lease.write();
if let Some(existing) = &*lease {
if existing.holder_pid == pid {
*lease = None;
drop(lease);
info!(bead_id = "bd-389e", pid, "coordinator lease released");
return true;
}
}
false
}
pub fn force_release_lease(&self) {
let mut lease = self.lease.write();
if let Some(existing) = &*lease {
warn!(
bead_id = "bd-389e",
pid = existing.holder_pid,
"coordinator lease force-released (crash recovery)"
);
}
*lease = None;
}
pub fn native_publish(&self, req: &NativePublishRequest) -> NativePublishResponse {
assert_eq!(
self.mode,
CoordinatorMode::Native,
"native_publish called in compatibility mode"
);
debug!(
bead_id = "bd-389e",
txn = ?req.txn,
pages = req.write_set_summary.len(),
"native_publish: starting validation"
);
if let Some((conflict_pages, conflict_seq)) =
self.validate_fcw_set(&req.write_set_summary, req.begin_seq)
{
info!(
bead_id = "bd-389e",
txn = ?req.txn,
conflicts = conflict_pages.len(),
"native_publish: FCW conflict detected"
);
return NativePublishResponse::Conflict {
conflicting_pages: conflict_pages,
conflicting_commit_seq: conflict_seq,
};
}
let commit_seq = self.allocate_commit_seq();
self.update_commit_index(&req.write_set_summary, commit_seq);
let marker_object_id = Self::derive_marker_id(req.txn, commit_seq);
info!(
bead_id = "bd-389e",
txn = ?req.txn,
commit_seq = commit_seq.get(),
"native_publish: commit approved (marker only, no page bytes)"
);
NativePublishResponse::Ok {
commit_seq,
marker_object_id,
}
}
pub fn compat_commit(&self, req: &CompatCommitRequest) -> CompatCommitResponse {
assert_eq!(
self.mode,
CoordinatorMode::Compatibility,
"compat_commit called in native mode"
);
let page_numbers: Vec<u32> = req
.write_set
.page_numbers()
.iter()
.map(|p| p.get())
.collect();
let page_set: BTreeSet<u32> = page_numbers.iter().copied().collect();
debug!(
bead_id = "bd-389e",
txn = ?req.txn,
mode = ?req.mode,
pages = page_numbers.len(),
spilled = req.write_set.is_spilled(),
"compat_commit: starting validation"
);
if let Some((conflict_pages, conflict_seq)) =
self.validate_fcw_set(&page_set, req.snapshot.high)
{
info!(
bead_id = "bd-389e",
txn = ?req.txn,
conflicts = conflict_pages.len(),
"compat_commit: FCW conflict detected"
);
return CompatCommitResponse::Conflict {
conflicting_pages: conflict_pages,
conflicting_commit_seq: conflict_seq,
};
}
let commit_seq = self.allocate_commit_seq();
let frame_header_size = 24_u64;
let page_size = Self::infer_page_size(&req.write_set);
let batch_bytes = page_numbers.len() as u64 * (frame_header_size + page_size);
let wal_offset = self.wal_offset.fetch_add(batch_bytes, Ordering::SeqCst);
self.update_commit_index(&page_set, commit_seq);
self.committed_seqs.write().push(commit_seq);
info!(
bead_id = "bd-389e",
txn = ?req.txn,
commit_seq = commit_seq.get(),
wal_offset,
pages = page_numbers.len(),
"compat_commit: commit approved (WAL path)"
);
CompatCommitResponse::Ok {
wal_offset,
commit_seq,
}
}
pub fn compat_commit_batch(
&self,
requests: &[CompatCommitRequest],
) -> Vec<CompatCommitResponse> {
assert_eq!(
self.mode,
CoordinatorMode::Compatibility,
"compat_commit_batch called in native mode"
);
let mut responses = Vec::with_capacity(requests.len());
let mut accepted_commits: Vec<(CommitSeq, BTreeSet<u32>)> = Vec::new();
let mut batch_page_owner: HashMap<u32, CommitSeq> = HashMap::new();
let frame_header_size = 24_u64;
let mut total_batch_bytes = 0_u64;
for req in requests {
let page_numbers: Vec<u32> = req
.write_set
.page_numbers()
.iter()
.map(|p| p.get())
.collect();
let page_set: BTreeSet<u32> = page_numbers.iter().copied().collect();
if let Some((conflict_pages, conflict_seq)) =
self.validate_fcw_set(&page_set, req.snapshot.high)
{
responses.push(CompatCommitResponse::Conflict {
conflicting_pages: conflict_pages,
conflicting_commit_seq: conflict_seq,
});
} else {
let mut intra_batch_conflicts = Vec::new();
let mut intra_batch_conflict_seq = CommitSeq::new(0);
for &pgno in &page_set {
if let Some(&owner_seq) = batch_page_owner.get(&pgno) {
if let Some(page) = PageNumber::new(pgno) {
intra_batch_conflicts.push(page);
}
if owner_seq.get() > intra_batch_conflict_seq.get() {
intra_batch_conflict_seq = owner_seq;
}
}
}
if intra_batch_conflicts.is_empty() {
let commit_seq = self.allocate_commit_seq();
let page_size = Self::infer_page_size(&req.write_set);
let page_count = req.write_set.page_count() as u64;
let commit_bytes = page_count * (frame_header_size + page_size);
let wal_offset = self.wal_offset.fetch_add(commit_bytes, Ordering::SeqCst);
total_batch_bytes += commit_bytes;
for &pgno in &page_set {
batch_page_owner.insert(pgno, commit_seq);
}
accepted_commits.push((commit_seq, page_set));
responses.push(CompatCommitResponse::Ok {
wal_offset,
commit_seq,
});
} else {
responses.push(CompatCommitResponse::Conflict {
conflicting_pages: intra_batch_conflicts,
conflicting_commit_seq: intra_batch_conflict_seq,
});
}
}
}
if accepted_commits.is_empty() {
return responses;
}
let accepted_count = accepted_commits.len();
debug!(
bead_id = "bd-389e",
batch_size = accepted_count,
total_bytes = total_batch_bytes,
"compat_commit_batch: single fsync for batch"
);
for (commit_seq, page_set) in &accepted_commits {
self.update_commit_index(page_set, *commit_seq);
self.committed_seqs.write().push(*commit_seq);
}
info!(
bead_id = "bd-389e",
batch_size = accepted_count,
conflicts = requests.len() - accepted_count,
"compat_commit_batch: group commit complete"
);
responses
}
fn validate_fcw_set(
&self,
write_pages: &BTreeSet<u32>,
begin_seq: CommitSeq,
) -> Option<(Vec<PageNumber>, CommitSeq)> {
let index = self.commit_page_index.read();
let mut conflict_pages = Vec::new();
let mut conflict_seq = CommitSeq::new(0);
for &pgno in write_pages {
if let Some(&committed_seq) = index.get(&pgno) {
if committed_seq.get() > begin_seq.get() {
if let Some(pn) = PageNumber::new(pgno) {
conflict_pages.push(pn);
}
if committed_seq.get() > conflict_seq.get() {
conflict_seq = committed_seq;
}
}
}
}
if conflict_pages.is_empty() {
None
} else {
Some((conflict_pages, conflict_seq))
}
}
fn allocate_commit_seq(&self) -> CommitSeq {
let seq = self.next_commit_seq.fetch_add(1, Ordering::SeqCst);
CommitSeq::new(seq)
}
fn update_commit_index(&self, pages: &BTreeSet<u32>, commit_seq: CommitSeq) {
let mut index = self.commit_page_index.write();
for &pgno in pages {
index.insert(pgno, commit_seq);
}
}
fn infer_page_size(write_set: &CommitWriteSet) -> u64 {
match write_set {
CommitWriteSet::Inline(pages) => {
pages.values().next().map_or(4096, |pd| pd.len() as u64)
}
CommitWriteSet::Spilled(spilled) => spilled
.pages
.values()
.next()
.map_or(4096, |loc| u64::from(loc.len)),
}
}
fn derive_marker_id(txn: TxnToken, commit_seq: CommitSeq) -> ObjectId {
let mut bytes = [0u8; 16];
bytes[..8].copy_from_slice(&txn.id.get().to_le_bytes());
bytes[8..16].copy_from_slice(&commit_seq.get().to_le_bytes());
ObjectId::from_bytes(bytes)
}
}
#[cfg(test)]
#[allow(clippy::too_many_lines)]
mod tests {
use super::*;
use fsqlite_types::{SchemaEpoch, TxnEpoch, TxnId};
fn test_token(id: u64) -> TxnToken {
TxnToken::new(TxnId::new(id).unwrap(), TxnEpoch::new(0))
}
fn test_snapshot(high: u64) -> Snapshot {
Snapshot {
high: CommitSeq::new(high),
schema_epoch: SchemaEpoch::new(1),
}
}
fn test_page_data(pgno: u32) -> PageData {
let mut data = vec![0u8; 4096];
data[..4].copy_from_slice(&pgno.to_le_bytes());
PageData::from_vec(data)
}
fn inline_write_set(pages: &[u32]) -> CommitWriteSet {
let mut map = HashMap::new();
for &pgno in pages {
map.insert(PageNumber::new(pgno).unwrap(), test_page_data(pgno));
}
CommitWriteSet::Inline(map)
}
#[test]
fn test_native_sequencer_tiny_marker() {
let coord = WriteCoordinator::new(CoordinatorMode::Native);
coord.acquire_lease(1, 0);
let req = NativePublishRequest {
txn: test_token(1),
begin_seq: CommitSeq::new(0),
capsule_object_id: ObjectId::from_bytes([1u8; 16]),
capsule_digest: [0xAB; 32],
write_set_summary: BTreeSet::from([5, 10, 15]),
read_witnesses: vec![ObjectId::from_bytes([2u8; 16])],
write_witnesses: vec![ObjectId::from_bytes([3u8; 16])],
edge_ids: Vec::new(),
merge_witnesses: Vec::new(),
abort_policy: AbortPolicy::AbortPivot,
};
let resp = coord.native_publish(&req);
match resp {
NativePublishResponse::Ok {
commit_seq,
marker_object_id,
} => {
assert!(commit_seq.get() > 0, "commit_seq must be positive");
assert_ne!(
marker_object_id,
ObjectId::from_bytes([0u8; 16]),
"marker must be non-zero"
);
}
other => panic!("expected Ok, got {other:?}"),
}
}
#[test]
fn test_compat_group_commit() {
let coord = WriteCoordinator::new(CoordinatorMode::Compatibility);
coord.acquire_lease(1, 0);
let requests: Vec<CompatCommitRequest> = (1..=3_u64)
.map(|i| {
#[allow(clippy::cast_possible_truncation)]
let pgno = (i as u32) * 10;
CompatCommitRequest {
txn: test_token(i),
mode: TransactionMode::Concurrent,
write_set: inline_write_set(&[pgno]),
intent_log: Vec::new(),
page_locks: HashSet::from([PageNumber::new(pgno).unwrap()]),
snapshot: test_snapshot(0),
has_in_rw: false,
has_out_rw: false,
wal_fec_r: 0,
}
})
.collect();
let responses = coord.compat_commit_batch(&requests);
assert_eq!(responses.len(), 3);
let mut commit_seqs = Vec::new();
for resp in &responses {
match resp {
CompatCommitResponse::Ok { commit_seq, .. } => {
commit_seqs.push(commit_seq.get());
}
other => panic!("expected Ok, got {other:?}"),
}
}
for window in commit_seqs.windows(2) {
assert!(window[0] < window[1], "commit_seqs must be monotonic");
}
}
#[test]
fn test_write_set_spill() {
let spill_loc = SpillLoc {
offset: 0,
len: 4096,
xxh3_64: 0xDEAD_BEEF,
};
let spill = SpilledWriteSet {
spill: SpillHandle::Path(PathBuf::from("/tmp/test-spill.dat")),
pages: HashMap::from([(PageNumber::new(5).unwrap(), spill_loc)]),
};
let write_set = CommitWriteSet::Spilled(spill);
assert!(write_set.is_spilled());
assert_eq!(write_set.page_count(), 1);
assert_eq!(write_set.page_numbers().len(), 1);
let coord = WriteCoordinator::new(CoordinatorMode::Compatibility);
coord.acquire_lease(1, 0);
let req = CompatCommitRequest {
txn: test_token(1),
mode: TransactionMode::Concurrent,
write_set,
intent_log: Vec::new(),
page_locks: HashSet::from([PageNumber::new(5).unwrap()]),
snapshot: test_snapshot(0),
has_in_rw: false,
has_out_rw: false,
wal_fec_r: 0,
};
let resp = coord.compat_commit(&req);
match resp {
CompatCommitResponse::Ok { commit_seq, .. } => {
assert!(commit_seq.get() > 0);
}
other => panic!("expected Ok, got {other:?}"),
}
}
#[test]
fn test_coordinator_lease() {
let coord = WriteCoordinator::new(CoordinatorMode::Native);
assert!(coord.acquire_lease(100, 0), "first acquire should succeed");
assert!(!coord.acquire_lease(200, 1), "second acquire should fail");
assert!(coord.release_lease(100), "release by holder should succeed");
assert!(
coord.acquire_lease(200, 2),
"acquire after release should succeed"
);
assert!(
!coord.release_lease(999),
"release by non-holder should fail"
);
}
#[test]
fn test_coordinator_role_takeover() {
let coord = WriteCoordinator::new(CoordinatorMode::Native);
assert!(coord.acquire_lease(100, 0));
coord.force_release_lease();
assert!(
coord.acquire_lease(200, 1),
"takeover after force-release should succeed"
);
}
#[test]
fn test_wal_frame_format() {
let coord = WriteCoordinator::new(CoordinatorMode::Compatibility);
coord.acquire_lease(1, 0);
let req = CompatCommitRequest {
txn: test_token(1),
mode: TransactionMode::Serialized,
write_set: inline_write_set(&[1, 2, 3]),
intent_log: Vec::new(),
page_locks: HashSet::from([
PageNumber::new(1).unwrap(),
PageNumber::new(2).unwrap(),
PageNumber::new(3).unwrap(),
]),
snapshot: test_snapshot(0),
has_in_rw: false,
has_out_rw: false,
wal_fec_r: 0,
};
let resp = coord.compat_commit(&req);
match resp {
CompatCommitResponse::Ok {
wal_offset,
commit_seq,
} => {
assert_eq!(wal_offset, 0, "first commit starts at WAL offset 0");
assert!(commit_seq.get() > 0);
let expected_next = 3 * (24 + 4096);
assert_eq!(
coord.wal_offset.load(Ordering::SeqCst),
expected_next,
"WAL offset advances by frame_header + page_size per page"
);
}
other => panic!("expected Ok, got {other:?}"),
}
let req2 = CompatCommitRequest {
txn: test_token(2),
mode: TransactionMode::Serialized,
write_set: inline_write_set(&[4]),
intent_log: Vec::new(),
page_locks: HashSet::from([PageNumber::new(4).unwrap()]),
snapshot: test_snapshot(0),
has_in_rw: false,
has_out_rw: false,
wal_fec_r: 0,
};
let resp2 = coord.compat_commit(&req2);
match resp2 {
CompatCommitResponse::Ok { wal_offset, .. } => {
assert_eq!(
wal_offset,
3 * (24 + 4096),
"second commit starts after first"
);
}
other => panic!("expected Ok, got {other:?}"),
}
}
#[test]
fn test_compat_group_commit_intra_batch_conflict_first_wins() {
let coord = WriteCoordinator::new(CoordinatorMode::Compatibility);
coord.acquire_lease(1, 0);
let requests = vec![
CompatCommitRequest {
txn: test_token(1),
mode: TransactionMode::Concurrent,
write_set: inline_write_set(&[42]),
intent_log: Vec::new(),
page_locks: HashSet::from([PageNumber::new(42).unwrap()]),
snapshot: test_snapshot(0),
has_in_rw: false,
has_out_rw: false,
wal_fec_r: 0,
},
CompatCommitRequest {
txn: test_token(2),
mode: TransactionMode::Concurrent,
write_set: inline_write_set(&[42]),
intent_log: Vec::new(),
page_locks: HashSet::from([PageNumber::new(42).unwrap()]),
snapshot: test_snapshot(0),
has_in_rw: false,
has_out_rw: false,
wal_fec_r: 0,
},
];
let responses = coord.compat_commit_batch(&requests);
assert_eq!(responses.len(), 2);
let first_commit_seq = match &responses[0] {
CompatCommitResponse::Ok { commit_seq, .. } => *commit_seq,
other => panic!("expected first response Ok, got {other:?}"),
};
match &responses[1] {
CompatCommitResponse::Conflict {
conflicting_pages,
conflicting_commit_seq,
} => {
assert_eq!(conflicting_pages, &vec![PageNumber::new(42).unwrap()]);
assert_eq!(*conflicting_commit_seq, first_commit_seq);
}
other => panic!("expected second response Conflict, got {other:?}"),
}
}
#[test]
fn test_restore_state_restores_wal_offset() {
let coord = WriteCoordinator::new(CoordinatorMode::Compatibility);
coord.acquire_lease(1, 0);
coord.restore_state(CommitSeq::new(11), HashMap::new(), 12_345);
let req = CompatCommitRequest {
txn: test_token(77),
mode: TransactionMode::Serialized,
write_set: inline_write_set(&[7]),
intent_log: Vec::new(),
page_locks: HashSet::from([PageNumber::new(7).unwrap()]),
snapshot: test_snapshot(0),
has_in_rw: false,
has_out_rw: false,
wal_fec_r: 0,
};
match coord.compat_commit(&req) {
CompatCommitResponse::Ok {
wal_offset,
commit_seq,
} => {
assert_eq!(wal_offset, 12_345);
assert_eq!(commit_seq, CommitSeq::new(11));
}
other => panic!("expected Ok, got {other:?}"),
}
}
}