use std::io;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::thread;
use std::time::Duration;
use kovan_queue::array_queue::ArrayQueue;
use super::wal::Wal;
use super::{CommitOutcome, DurabilityMode, RegolithEngine, grouped_batch_ops};
use crate::WriteBatchOp;
use crate::perf_context::{PerfTimer, PerfTimerField};
use crate::statistics::{Histogram, Ticker};
mod request;
mod slot;
mod stall;
pub(crate) use request::WriteRequest;
pub(crate) use slot::WriteSlot;
pub(crate) use stall::StallSignal;
const MAX_GROUP_BYTES: usize = 1024 * 1024;
const PARK_SLICE: Duration = Duration::from_micros(200);
fn commit_ring_capacity() -> usize {
thread::available_parallelism()
.map(|n| n.get().saturating_mul(4))
.unwrap_or(16)
.clamp(16, 1024)
}
struct GroupTicket {
slot: Option<Arc<WriteSlot>>,
request: WriteRequest,
}
pub(crate) struct Pipeline {
stage: Vec<u8>,
group: Vec<GroupTicket>,
}
impl Pipeline {
pub(crate) fn new() -> Self {
Self {
stage: Vec::new(),
group: Vec::new(),
}
}
}
fn clone_io_error(err: &io::Error) -> io::Error {
io::Error::new(err.kind(), err.to_string())
}
fn release_stranded(group: &mut Vec<GroupTicket>) {
for ticket in group.drain(..) {
if let Some(slot) = ticket.slot {
slot.complete(Err(io::Error::other(
"commit group abandoned by a leader that did not finish",
)));
}
}
}
impl RegolithEngine {
pub(crate) fn new_commit_ring() -> ArrayQueue<Arc<WriteSlot>> {
ArrayQueue::new(commit_ring_capacity())
}
pub(crate) fn apply_batch(
&self,
ops: Vec<WriteBatchOp>,
durability: DurabilityMode,
disable_wal: bool,
) -> io::Result<u64> {
self.ensure_writable()?;
if ops.is_empty() {
return Ok(self.visible_seq.visible());
}
self.validate_ops_sizes(&ops)?;
self.submit(WriteRequest::Batch {
ops,
durability,
disable_wal,
})
}
pub(crate) fn apply_grouped_batch(
&self,
point_ops: std::collections::BTreeMap<Vec<u8>, Option<Vec<u8>>>,
range_deletes: Vec<(Vec<u8>, Vec<u8>)>,
merges: Vec<(Vec<u8>, Vec<u8>)>,
durability: DurabilityMode,
disable_wal: bool,
) -> io::Result<u64> {
let ops = grouped_batch_ops(point_ops, range_deletes, merges);
self.apply_batch(ops, durability, disable_wal)
}
pub(crate) fn apply_single_put(
&self,
key: Vec<u8>,
value: Vec<u8>,
durability: DurabilityMode,
disable_wal: bool,
) -> io::Result<u64> {
self.ensure_writable()?;
self.validate_prefixed_key_size(&key)?;
self.validate_value_size(&value)?;
self.submit(WriteRequest::Put {
key,
value,
durability,
disable_wal,
})
}
pub(crate) fn commit_optimistic(
&self,
conflict_keys: &[(Vec<u8>, u64)],
point_ops: std::collections::BTreeMap<Vec<u8>, Option<Vec<u8>>>,
range_deletes: Vec<(Vec<u8>, Vec<u8>)>,
merges: Vec<(Vec<u8>, Vec<u8>)>,
durability: DurabilityMode,
) -> io::Result<CommitOutcome> {
self.ensure_writable()?;
let ops = grouped_batch_ops(point_ops, range_deletes, merges);
self.validate_ops_sizes(&ops)?;
let mut pipe = self.pipeline.lock();
let view = self.view.load();
for (key, observed_seq) in conflict_keys {
if let Some(latest_seq) = self.latest_version_seq_in_view(key, &view)?
&& latest_seq > *observed_seq
{
return Ok(CommitOutcome::Conflict {
key: key.clone(),
observed_seq: *observed_seq,
latest_seq,
});
}
}
if ops.is_empty() {
return Ok(CommitOutcome::Ok);
}
let request = WriteRequest::Batch {
ops,
durability,
disable_wal: false,
};
release_stranded(&mut pipe.group);
pipe.group.push(GroupTicket {
slot: None,
request,
});
let result = self.run_and_complete(&mut pipe);
self.drain_locked(&mut pipe);
result.map(|_| CommitOutcome::Ok)
}
fn submit(&self, request: WriteRequest) -> io::Result<u64> {
if let Some(mut pipe) = self.pipeline.try_lock() {
return self.lead_with(&mut pipe, request);
}
let slot = slot::thread_slot();
if let Err(request) = slot.arm(request) {
let mut pipe = self.pipeline.lock();
return self.lead_with(&mut pipe, request);
}
let mut ticket = Arc::clone(&slot);
while let Err(returned) = self.commit_ring.push(ticket) {
ticket = returned;
if !self.try_drain() {
thread::yield_now();
}
}
while !slot.is_done() {
if self.try_drain() {
continue;
}
thread::park_timeout(PARK_SLICE);
}
slot.finish()
}
pub(super) fn lead_with(&self, pipe: &mut Pipeline, request: WriteRequest) -> io::Result<u64> {
release_stranded(&mut pipe.group);
pipe.group.push(GroupTicket {
slot: None,
request,
});
self.admit_from_ring(pipe);
let result = self.run_and_complete(pipe);
self.drain_locked(pipe);
result
}
fn try_drain(&self) -> bool {
let Some(mut pipe) = self.pipeline.try_lock() else {
return false;
};
self.drain_locked(&mut pipe);
true
}
fn drain_locked(&self, pipe: &mut Pipeline) {
loop {
release_stranded(&mut pipe.group);
self.admit_from_ring(pipe);
if pipe.group.is_empty() {
return;
}
let _ = self.run_and_complete(pipe);
}
}
fn admit_from_ring(&self, pipe: &mut Pipeline) {
let room = self.memtable_room();
let mut staged: usize = pipe.group.iter().map(|t| t.request.staged_len()).sum();
let mut projected: usize = pipe.group.iter().map(|t| t.request.memtable_cost()).sum();
loop {
if !pipe.group.is_empty() && (staged >= MAX_GROUP_BYTES || projected >= room) {
return;
}
let Some(slot) = self.commit_ring.pop() else {
return;
};
let request = slot.take_request();
staged += request.staged_len();
projected += request.memtable_cost();
pipe.group.push(GroupTicket {
slot: Some(slot),
request,
});
}
}
fn memtable_room(&self) -> usize {
let budget = self.options.write_buffer_size;
let used = self.view.load().active.approximate_size();
if used >= budget {
budget
} else {
budget - used
}
}
fn run_and_complete(&self, pipe: &mut Pipeline) -> io::Result<u64> {
let Pipeline { stage, group } = pipe;
let result = self.run_group(stage, group);
let mut seq = result.as_ref().ok().copied().unwrap_or(0);
for ticket in group.drain(..) {
let ops = ticket.request.op_count();
let last = seq.saturating_add(ops).saturating_sub(1);
if let Some(slot) = ticket.slot {
slot.complete(match &result {
Ok(_) => Ok(last),
Err(e) => Err(io::Error::new(e.kind(), e.to_string())),
});
}
seq = last.saturating_add(1);
}
if stage.capacity() > MAX_GROUP_BYTES {
stage.clear();
stage.shrink_to(MAX_GROUP_BYTES);
}
result
}
fn run_group(&self, stage: &mut Vec<u8>, group: &[GroupTicket]) -> io::Result<u64> {
self.ensure_writable()?;
self.rotate_if_full()?;
let total_ops: u64 = group.iter().map(|t| t.request.op_count()).sum();
if total_ops == 0 {
return Ok(self.visible_seq.visible());
}
let base_seq = self.latest_seq.fetch_add(total_ops, Ordering::AcqRel) + 1;
stage.clear();
let mut any_immediate = false;
let mut reported_bytes = 0u64;
let mut seq = base_seq;
for ticket in group {
if !ticket.request.skips_wal() {
ticket.request.encode_wal(stage, seq);
reported_bytes += ticket.request.reported_bytes();
any_immediate |= matches!(ticket.request.durability(), DurabilityMode::Immediate);
}
seq += ticket.request.op_count();
}
if !stage.is_empty() {
let _perf_wal = PerfTimer::new(PerfTimerField::WriteWal);
let wal_start = self.statistics().and_then(|_| self.env.now_micros());
let mut guard = self.active_wal.lock();
let wal = guard.as_mut().ok_or_else(Self::read_only_error)?;
let start_offset = wal.offset();
if let Err(err) = wal.append_group(stage) {
self.abandon_group(wal, start_offset, &err);
return Err(err);
}
let mut synced = 0u64;
if any_immediate {
if let Err(err) = wal.sync_data() {
self.abandon_group(wal, start_offset, &err);
return Err(err);
}
synced = 1;
}
drop(guard);
if let Some(s) = self.statistics() {
s.add(Ticker::WalBytesWritten, reported_bytes);
if synced > 0 {
s.add(Ticker::WalSyncCount, synced);
}
if let Some(micros) = self.elapsed_micros(wal_start) {
s.record(Histogram::WalWriteTime, micros);
}
}
}
{
let _perf_mt = PerfTimer::new(PerfTimerField::WriteMemtable);
let view = self.view.load();
let memtable = &view.active;
let mut seq = base_seq;
for ticket in group {
ticket.request.apply(memtable, &mut seq);
}
}
self.visible_seq.publish(base_seq + total_ops - 1);
Ok(base_seq)
}
fn abandon_group(&self, wal: &mut Wal, start_offset: u64, cause: &io::Error) {
tracing::error!(error = %cause, "commit group failed; discarding its WAL bytes");
if let Err(rollback_err) = wal.rollback_to(start_offset) {
self.latch_wal_failure(&rollback_err);
}
if !self.options.listeners.is_empty() {
let err = crate::Error::from(clone_io_error(cause));
crate::event_listener::dispatch(&self.options.listeners, |l| {
l.on_background_error(
crate::event_listener::BackgroundErrorReason::WriteAheadLog,
&err,
)
});
}
}
}
#[cfg(test)]
mod tests {
use super::super::{EngineOptions, wal::fault};
use super::*;
use crate::sync::Mutex;
use tempfile::TempDir;
fn open_engine(dir: &TempDir) -> Arc<RegolithEngine> {
RegolithEngine::open(dir.path(), EngineOptions::default()).expect("engine open")
}
fn key(name: &[u8]) -> Vec<u8> {
let mut k = vec![0u8; 4];
k.extend_from_slice(name);
k
}
fn durable_put(name: &[u8], value: &[u8]) -> WriteRequest {
WriteRequest::Put {
key: key(name),
value: value.to_vec(),
durability: DurabilityMode::Immediate,
disable_wal: false,
}
}
struct FaultGuard(std::path::PathBuf);
impl Drop for FaultGuard {
fn drop(&mut self) {
fault::disarm_sync_failure(&self.0);
}
}
fn arm_sync_failure(dir: &TempDir) -> FaultGuard {
fault::arm_sync_failure(dir.path());
FaultGuard(dir.path().to_path_buf())
}
fn arm_flapping_sync_failure(dir: &TempDir) -> FaultGuard {
fault::arm_flapping_sync_failure(dir.path(), 2);
FaultGuard(dir.path().to_path_buf())
}
#[test]
fn cloned_errors_keep_kind_and_message() {
let err = io::Error::new(io::ErrorKind::StorageFull, "disk is full");
let cloned = clone_io_error(&err);
assert_eq!(cloned.kind(), err.kind());
assert_eq!(cloned.to_string(), err.to_string());
}
#[test]
fn commit_ring_capacity_is_bounded() {
let cap = commit_ring_capacity();
assert!((16..=1024).contains(&cap), "unexpected ring capacity {cap}");
}
#[test]
fn a_failed_sync_fails_every_member_of_the_group() {
let dir = TempDir::new().unwrap();
let engine = open_engine(&dir);
let horizon_before = engine.snapshot_seq();
let wal_len_before = engine
.active_wal
.lock()
.as_ref()
.map(|w| w.offset())
.expect("writable engine has a wal");
let followers: Vec<Arc<WriteSlot>> = (0..3)
.map(|i| {
let slot = Arc::new(WriteSlot::new());
let request = durable_put(format!("member{i}").as_bytes(), b"value");
slot.arm(request).expect("fresh slot arms");
slot
})
.collect();
{
let _fault = arm_sync_failure(&dir);
let mut pipe = engine.pipeline.lock();
pipe.group.clear();
for slot in &followers {
let request = slot.take_request();
pipe.group.push(GroupTicket {
slot: Some(Arc::clone(slot)),
request,
});
}
let result = engine.run_and_complete(&mut pipe);
assert!(
result.is_err(),
"an injected sync failure must fail the group"
);
}
for (i, slot) in followers.iter().enumerate() {
assert!(slot.is_done(), "member {i} was left pending");
let outcome = slot.finish();
assert!(
outcome.is_err(),
"member {i} must learn the group did not commit"
);
}
assert_eq!(
engine.snapshot_seq(),
horizon_before,
"a failed group must not publish a read horizon"
);
for i in 0..3 {
let name = format!("member{i}");
assert_eq!(
engine.get(&key(name.as_bytes()), u64::MAX).unwrap(),
None,
"a failed group must not be applied to the memtable"
);
}
assert_eq!(
engine
.active_wal
.lock()
.as_ref()
.map(|w| w.offset())
.unwrap(),
wal_len_before,
"a failed group must be rolled back out of the WAL"
);
}
#[test]
fn a_failed_group_leaves_nothing_to_recover() {
let dir = TempDir::new().unwrap();
{
let engine = open_engine(&dir);
engine
.submit(durable_put(b"before", b"kept"))
.expect("the pre-failure write commits");
let _fault = arm_sync_failure(&dir);
let err = engine
.submit(durable_put(b"during", b"lost"))
.expect_err("the injected failure must surface");
assert!(err.to_string().contains("injected"));
}
let engine = open_engine(&dir);
assert_eq!(
engine.get(&key(b"before"), u64::MAX).unwrap(),
Some(b"kept".to_vec())
);
assert_eq!(
engine.get(&key(b"during"), u64::MAX).unwrap(),
None,
"a write whose group failed must not survive a reopen"
);
}
#[test]
fn an_abandoned_ticket_is_completed_and_does_not_wedge_the_ring() {
let dir = TempDir::new().unwrap();
let engine = open_engine(&dir);
let orphan = Arc::new(WriteSlot::new());
orphan
.arm(durable_put(b"orphan", b"value"))
.expect("fresh slot arms");
engine
.commit_ring
.push(Arc::clone(&orphan))
.map_err(|_| "commit ring full")
.expect("empty ring accepts a ticket");
engine
.submit(durable_put(b"later", b"value"))
.expect("a later writer drains the ring");
assert!(orphan.is_done(), "the abandoned ticket was never executed");
assert!(orphan.finish().is_ok());
assert_eq!(
engine.get(&key(b"orphan"), u64::MAX).unwrap(),
Some(b"value".to_vec())
);
assert_eq!(
engine.get(&key(b"later"), u64::MAX).unwrap(),
Some(b"value".to_vec())
);
}
#[test]
fn one_group_costs_one_sync_no_matter_how_many_members() {
let dir = TempDir::new().unwrap();
let stats = Arc::new(crate::statistics::Statistics::new());
let engine = RegolithEngine::open(
dir.path(),
EngineOptions {
statistics: Some(Arc::clone(&stats)),
..EngineOptions::default()
},
)
.unwrap();
let followers: Vec<Arc<WriteSlot>> = (0..8)
.map(|i| {
let slot = Arc::new(WriteSlot::new());
slot.arm(durable_put(format!("k{i}").as_bytes(), b"v"))
.expect("fresh slot arms");
engine
.commit_ring
.push(Arc::clone(&slot))
.map_err(|_| "commit ring full")
.expect("ring accepts the ticket");
slot
})
.collect();
assert_eq!(stats.get_ticker(Ticker::WalSyncCount), 0);
engine.try_drain();
for slot in &followers {
assert!(slot.is_done());
slot.finish().expect("every member commits");
}
assert_eq!(
stats.get_ticker(Ticker::WalSyncCount),
1,
"eight durable writers in one group must cost exactly one fdatasync"
);
for i in 0..8 {
let name = format!("k{i}");
assert_eq!(
engine.get(&key(name.as_bytes()), u64::MAX).unwrap(),
Some(b"v".to_vec())
);
}
}
#[test]
fn a_group_publishes_the_horizon_only_after_every_member_is_applied() {
let dir = TempDir::new().unwrap();
let engine = open_engine(&dir);
let slots: Vec<Arc<WriteSlot>> = (0..4)
.map(|i| {
let slot = Arc::new(WriteSlot::new());
slot.arm(durable_put(format!("h{i}").as_bytes(), b"v"))
.expect("fresh slot arms");
slot
})
.collect();
let mut pipe = engine.pipeline.lock();
pipe.group.clear();
for slot in &slots {
let request = slot.take_request();
pipe.group.push(GroupTicket {
slot: Some(Arc::clone(slot)),
request,
});
}
engine.run_and_complete(&mut pipe).expect("group commits");
drop(pipe);
let horizon = engine.snapshot_seq();
for i in 0..4 {
let name = format!("h{i}");
assert_eq!(
engine.get(&key(name.as_bytes()), horizon).unwrap(),
Some(b"v".to_vec()),
"the published horizon must cover every member of the group"
);
}
}
#[test]
fn a_group_stranded_by_an_unwound_leader_is_released_not_left_parked() {
let dir = TempDir::new().unwrap();
let engine = open_engine(&dir);
let stranded = Arc::new(WriteSlot::new());
stranded
.arm(durable_put(b"stranded", b"v"))
.expect("fresh slot arms");
{
let mut pipe = engine.pipeline.lock();
let request = stranded.take_request();
pipe.group.push(GroupTicket {
slot: Some(Arc::clone(&stranded)),
request,
});
}
engine
.submit(durable_put(b"next", b"v"))
.expect("the next leader commits its own write");
assert!(stranded.is_done(), "the stranded writer is still parked");
let err = stranded
.finish()
.expect_err("a stranded writer must not be told it committed");
assert!(err.to_string().contains("abandoned"));
assert_eq!(
engine.get(&key(b"next"), u64::MAX).unwrap(),
Some(b"v".to_vec())
);
}
#[test]
fn an_outsized_request_does_not_park_its_staging_buffer() {
let dir = TempDir::new().unwrap();
let engine = open_engine(&dir);
let huge = vec![b'x'; MAX_GROUP_BYTES + 64 * 1024];
engine
.submit(WriteRequest::Put {
key: key(b"huge"),
value: huge.clone(),
durability: DurabilityMode::Eventual,
disable_wal: false,
})
.expect("a request larger than the group cap still commits");
assert!(
engine.pipeline.lock().stage.capacity() <= MAX_GROUP_BYTES,
"the staging buffer kept an outsized request's peak allocation"
);
assert_eq!(engine.get(&key(b"huge"), u64::MAX).unwrap(), Some(huge));
}
#[test]
fn a_wal_disabled_member_rides_along_without_forcing_a_sync() {
let dir = TempDir::new().unwrap();
let stats = Arc::new(crate::statistics::Statistics::new());
let engine = RegolithEngine::open(
dir.path(),
EngineOptions {
statistics: Some(Arc::clone(&stats)),
..EngineOptions::default()
},
)
.unwrap();
let quiet = Arc::new(WriteSlot::new());
quiet
.arm(WriteRequest::Put {
key: key(b"quiet"),
value: b"v".to_vec(),
durability: DurabilityMode::Eventual,
disable_wal: true,
})
.expect("fresh slot arms");
engine
.commit_ring
.push(Arc::clone(&quiet))
.map_err(|_| "commit ring full")
.expect("ring accepts the ticket");
engine
.submit(WriteRequest::Put {
key: key(b"eventual"),
value: b"v".to_vec(),
durability: DurabilityMode::Eventual,
disable_wal: false,
})
.expect("the group commits");
assert!(quiet.is_done());
quiet.finish().expect("the wal-disabled member commits");
assert_eq!(
stats.get_ticker(Ticker::WalSyncCount),
0,
"no member asked for Immediate durability, so no fsync is due"
);
assert_eq!(
engine.get(&key(b"quiet"), u64::MAX).unwrap(),
Some(b"v".to_vec())
);
}
#[test]
fn concurrent_writers_share_far_fewer_fsyncs_than_writes() {
let dir = TempDir::new().unwrap();
let stats = Arc::new(crate::statistics::Statistics::new());
let engine = RegolithEngine::open(
dir.path(),
EngineOptions {
statistics: Some(Arc::clone(&stats)),
..EngineOptions::default()
},
)
.unwrap();
const WRITERS: usize = 8;
const PER_WRITER: usize = 64;
let mut handles = Vec::with_capacity(WRITERS);
for w in 0..WRITERS {
let engine = Arc::clone(&engine);
handles.push(thread::spawn(move || {
for i in 0..PER_WRITER {
engine
.submit(durable_put(format!("w{w}k{i:04}").as_bytes(), b"value"))
.expect("every durable write commits");
}
}));
}
for handle in handles {
handle.join().expect("writer thread panicked");
}
let total = (WRITERS * PER_WRITER) as u64;
let syncs = stats.get_ticker(Ticker::WalSyncCount);
assert!(syncs >= 1, "durable writes must fsync at least once");
assert!(
syncs <= total,
"group commit can never issue more fsyncs than writes: {syncs} > {total}"
);
for w in 0..WRITERS {
for i in 0..PER_WRITER {
let name = format!("w{w}k{i:04}");
assert_eq!(
engine.get(&key(name.as_bytes()), u64::MAX).unwrap(),
Some(b"value".to_vec()),
"every committed write must be readable"
);
}
}
}
#[test]
fn real_concurrent_writers_all_learn_a_failed_fsync() {
let dir = TempDir::new().unwrap();
const WRITERS: usize = 12;
const PER_WRITER: usize = 40;
let acknowledged = {
let engine = open_engine(&dir);
engine
.submit(durable_put(b"before", b"kept"))
.expect("the pre-failure write commits");
let _fault = arm_flapping_sync_failure(&dir);
let acknowledged = Arc::new(Mutex::new(Vec::new()));
let mut handles = Vec::with_capacity(WRITERS);
for w in 0..WRITERS {
let engine = Arc::clone(&engine);
let acknowledged = Arc::clone(&acknowledged);
handles.push(thread::spawn(move || {
let durability = if w % 2 == 0 {
DurabilityMode::Immediate
} else {
DurabilityMode::Eventual
};
for i in 0..PER_WRITER {
let name = format!("f{w:02}_{i:03}");
let request = WriteRequest::Put {
key: key(name.as_bytes()),
value: b"v".to_vec(),
durability,
disable_wal: false,
};
if engine.submit(request).is_ok() {
acknowledged.lock().push(name);
}
}
}));
}
for handle in handles {
handle.join().expect("writer thread panicked");
}
let acknowledged = acknowledged.lock().clone();
assert!(
!acknowledged.is_empty() && acknowledged.len() < WRITERS * PER_WRITER,
"the flapping fault must produce a mix: {} of {} acknowledged",
acknowledged.len(),
WRITERS * PER_WRITER
);
for w in 0..WRITERS {
for i in 0..PER_WRITER {
let name = format!("f{w:02}_{i:03}");
let present = engine
.get(&key(name.as_bytes()), u64::MAX)
.unwrap()
.is_some();
assert_eq!(
present,
acknowledged.contains(&name),
"{name}: readable={present} but acknowledged={}",
acknowledged.contains(&name)
);
}
}
acknowledged
};
let engine = open_engine(&dir);
assert_eq!(
engine.get(&key(b"before"), u64::MAX).unwrap(),
Some(b"kept".to_vec()),
"a write that committed before the fault must survive"
);
for w in 0..WRITERS {
for i in 0..PER_WRITER {
let name = format!("f{w:02}_{i:03}");
let recovered = engine.get(&key(name.as_bytes()), u64::MAX).unwrap();
if acknowledged.contains(&name) {
assert_eq!(
recovered,
Some(b"v".to_vec()),
"{name} was acknowledged but lost across a reopen"
);
} else {
assert_eq!(
recovered, None,
"{name} was rejected but resurrected across a reopen"
);
}
}
}
}
#[test]
fn a_failed_group_does_not_truncate_an_earlier_groups_bytes() {
let dir = TempDir::new().unwrap();
{
let engine = open_engine(&dir);
engine
.submit(WriteRequest::Put {
key: key(b"eventual"),
value: b"kept".to_vec(),
durability: DurabilityMode::Eventual,
disable_wal: false,
})
.expect("the eventual write commits");
let _fault = arm_sync_failure(&dir);
engine
.submit(durable_put(b"doomed", b"lost"))
.expect_err("the injected failure must surface");
}
let engine = open_engine(&dir);
assert_eq!(
engine.get(&key(b"eventual"), u64::MAX).unwrap(),
Some(b"kept".to_vec()),
"a failed group rolled back past an earlier group's bytes"
);
assert_eq!(engine.get(&key(b"doomed"), u64::MAX).unwrap(), None);
}
#[test]
fn a_writer_that_commits_is_immediately_visible_to_another_thread() {
let dir = TempDir::new().unwrap();
let engine = open_engine(&dir);
let (tx, rx) = std::sync::mpsc::channel::<usize>();
let writer = {
let engine = Arc::clone(&engine);
thread::spawn(move || {
for i in 0..2_000usize {
engine
.submit(durable_put(format!("v{i:05}").as_bytes(), b"v"))
.expect("write commits");
tx.send(i).expect("reader is alive");
}
})
};
let mut checked = 0usize;
for i in rx {
let horizon = engine.snapshot_seq();
assert_eq!(
engine
.get(&key(format!("v{i:05}").as_bytes()), horizon)
.unwrap(),
Some(b"v".to_vec()),
"a committed write was not visible at the published horizon"
);
checked += 1;
}
writer.join().expect("writer thread panicked");
assert_eq!(checked, 2_000);
}
#[test]
fn writers_make_progress_with_more_threads_than_ring_slots() {
let dir = TempDir::new().unwrap();
let engine = open_engine(&dir);
let writers = commit_ring_capacity() + 8;
let mut handles = Vec::with_capacity(writers);
for w in 0..writers {
let engine = Arc::clone(&engine);
handles.push(thread::spawn(move || {
engine
.submit(durable_put(format!("overflow{w:04}").as_bytes(), b"v"))
.expect("a full ring must not lose a write");
}));
}
for handle in handles {
handle.join().expect("writer thread panicked");
}
for w in 0..writers {
let name = format!("overflow{w:04}");
assert_eq!(
engine.get(&key(name.as_bytes()), u64::MAX).unwrap(),
Some(b"v".to_vec())
);
}
}
}