use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use smallvec::SmallVec;
use tempfile::TempDir;
use tephra::Position;
use tephra::event::{Event, EventRef, EventType, Tag, Tags};
use tephra::log::set::{SegmentConfig, SegmentSet};
use tephra::query::{Matches, Query, QueryItem};
use tephra::read::{ReadConfig, ReadError};
use tephra::writer::{WriteCoordinator, WriteHandle, WriterConfig};
struct Rng(u64);
impl Rng {
fn next(&mut self) -> u64 {
self.0 = self
.0
.wrapping_mul(6364136223846793005)
.wrapping_add(1442695040888963407);
self.0 >> 17
}
fn below(&mut self, n: u64) -> u64 {
self.next() % n
}
}
const TYPES: [&str; 3] = ["Registered", "Enrolled", "Renamed"];
const TAGS: [&str; 6] = [
"course:c1",
"course:c2",
"student:s1",
"student:s2",
"team:t1",
"team:t2",
];
fn event_type(s: &str) -> EventType {
EventType::new(s).unwrap()
}
fn pick_tags(start: usize, k: usize) -> Tags {
let picked: SmallVec<[Tag; 4]> = (0..k)
.map(|i| Tag::new(TAGS[(start + i) % TAGS.len()]).unwrap())
.collect();
Tags::new(picked).unwrap()
}
fn random_event(rng: &mut Rng) -> Event {
let ty = event_type(TYPES[rng.below(TYPES.len() as u64) as usize]);
let k = rng.below(4) as usize;
let start = rng.below(TAGS.len() as u64) as usize;
Event::new(&ty, &pick_tags(start, k), b"payload").unwrap()
}
fn random_query(rng: &mut Rng) -> Query {
match rng.below(6) {
0 => Query::all(),
1 => Query::items(Vec::new()),
_ => {
let n_items = 1 + rng.below(3) as usize;
let items = (0..n_items)
.map(|_| {
let n_types = rng.below(3) as usize;
let type_start = rng.below(TYPES.len() as u64) as usize;
let types = (0..n_types)
.map(|i| event_type(TYPES[(type_start + i) % TYPES.len()]))
.collect();
let n_tags = rng.below(4) as usize;
let tag_start = rng.below(TAGS.len() as u64) as usize;
QueryItem::new(types, pick_tags(tag_start, n_tags))
})
.collect::<Vec<_>>();
Query::items(items)
}
}
}
fn scan_baseline(set: &SegmentSet, query: &Query, after: Position) -> Vec<Position> {
let mut out = Vec::new();
let mut scan = set.scan_after(after);
while let Some(item) = scan.next() {
let record = item.unwrap();
let event = EventRef::from_bytes(record.data).unwrap();
if query.matches(event) {
out.push(record.position);
}
}
out
}
fn coordinator() -> (WriteCoordinator, WriteHandle, TempDir) {
let dir = TempDir::new().unwrap();
let set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
let cfg = WriterConfig {
queue_capacity: 64,
max_batch_records: 64,
max_batch_bytes: 256,
tips_window: 1_000_000,
verify_tips: false,
..WriterConfig::default()
};
let (coord, handle) = WriteCoordinator::start(set, cfg).unwrap();
(coord, handle, dir)
}
fn coordinator_with_scan_bias(scan_bias: u32) -> (WriteCoordinator, WriteHandle, TempDir) {
let dir = TempDir::new().unwrap();
let set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
let cfg = WriterConfig {
queue_capacity: 64,
max_batch_records: 64,
max_batch_bytes: 256,
tips_window: 1_000_000,
verify_tips: false,
read: ReadConfig { scan_bias },
..WriterConfig::default()
};
let (coord, handle) = WriteCoordinator::start(set, cfg).unwrap();
(coord, handle, dir)
}
fn read_owned(
handle: &WriteHandle,
query: &Query,
after: Position,
) -> Result<Vec<(Position, Event)>, ReadError> {
handle.read(query, after, None).collect_owned()
}
fn tags(items: &[&str]) -> Tags {
Tags::new(
items
.iter()
.map(|s| Tag::new(*s).unwrap())
.collect::<SmallVec<[Tag; 4]>>(),
)
.unwrap()
}
fn tagged_event(ty: &str, tag_strs: &[&str]) -> Event {
Event::new(&EventType::new(ty).unwrap(), &tags(tag_strs), b"payload").unwrap()
}
#[test]
fn read_matches_scan_over_random_workload_across_segments() {
let (coord, handle, _dir) = coordinator();
let mut rng = Rng(0x1234_5678_9ABC_DEF0);
for _ in 0..400 {
let event = random_event(&mut rng);
handle.append(vec![event], None).unwrap();
}
type ReadCase = (Query, Position, Vec<(Position, Event)>);
let mut qrng = Rng(0xC0FF_EE00_1234_5678);
let last = 400u64;
let mut cases: Vec<ReadCase> = Vec::new();
for _ in 0..2000 {
let query = random_query(&mut qrng);
let after = Position::new(qrng.below(last + 1));
let got = read_owned(&handle, &query, after).unwrap();
cases.push((query, after, got));
}
let set = coord.shutdown();
for (query, after, got) in &cases {
let positions: Vec<Position> = got.iter().map(|(p, _)| *p).collect();
let baseline = scan_baseline(&set, query, *after);
assert_eq!(
positions, baseline,
"read positions disagreed with scan for query {query:?} after {after}"
);
for (position, event) in got {
let record = set.read_at(*position).unwrap();
let expected = EventRef::from_bytes(&record.data).unwrap();
assert_eq!(event.as_ref().event_type(), expected.event_type());
assert_eq!(event.as_ref().data(), expected.data());
let got_tags: Vec<&str> = event.as_ref().tags().collect();
let want_tags: Vec<&str> = expected.tags().collect();
assert_eq!(got_tags, want_tags);
}
}
}
#[test]
fn the_planner_never_changes_the_answer() {
let mut wrng = Rng(0x5EED_1234_ABCD_0001);
let events: Vec<Event> = (0..400).map(|_| random_event(&mut wrng)).collect();
let mut qrng = Rng(0xABCD_0001_5EED_1234);
let last = events.len() as u64;
let cases: Vec<(Query, Position)> = (0..500)
.map(|_| {
let query = random_query(&mut qrng);
let after = Position::new(qrng.below(last + 1));
(query, after)
})
.collect();
let mut per_bias: Vec<Vec<Vec<Position>>> = Vec::new();
for &scan_bias in &[1u32, 4, u32::MAX] {
let (coord, handle, _dir) = coordinator_with_scan_bias(scan_bias);
for ev in &events {
handle.append(vec![ev.clone()], None).unwrap();
}
let results: Vec<Vec<Position>> = cases
.iter()
.map(|(query, after)| {
read_owned(&handle, query, *after)
.unwrap()
.into_iter()
.map(|(p, _)| p)
.collect()
})
.collect();
let set = coord.shutdown();
if per_bias.is_empty() {
let baseline: Vec<Vec<Position>> = cases
.iter()
.map(|(query, after)| scan_baseline(&set, query, *after))
.collect();
assert_eq!(
results, baseline,
"scan_bias {scan_bias} disagreed with the scan oracle"
);
}
per_bias.push(results);
}
for (bias, results) in [4u32, u32::MAX].iter().zip(&per_bias[1..]) {
assert_eq!(
*results, per_bias[0],
"scan_bias {bias} changed the answer versus the forced-index run"
);
}
}
#[test]
fn a_limit_truncates_the_result_identically_on_every_path() {
let mut wrng = Rng(0x0FF1_CE00_1234_5678);
let events: Vec<Event> = (0..400).map(|_| random_event(&mut wrng)).collect();
let last = events.len() as u64;
let mut qrng = Rng(0x1234_5678_0FF1_CE00);
let cases: Vec<(Query, Position, u64)> = (0..500)
.map(|_| {
let query = random_query(&mut qrng);
let after = Position::new(qrng.below(last + 1));
let limit = qrng.below(6) * qrng.below(40); (query, after, limit)
})
.collect();
for &scan_bias in &[1u32, u32::MAX] {
let (coord, handle, _dir) = coordinator_with_scan_bias(scan_bias);
for ev in &events {
handle.append(vec![ev.clone()], None).unwrap();
}
let limited: Vec<Vec<Position>> = cases
.iter()
.map(|(query, after, limit)| {
handle
.read(query, *after, Some(*limit))
.collect_owned()
.unwrap()
.into_iter()
.map(|(p, _)| p)
.collect()
})
.collect();
let set = coord.shutdown();
for ((query, after, limit), got) in cases.iter().zip(&limited) {
let mut want = scan_baseline(&set, query, *after);
want.truncate(*limit as usize);
assert_eq!(
*got, want,
"scan_bias {scan_bias}: limited read disagreed with the truncated oracle \
(after {after:?}, limit {limit})"
);
}
}
}
#[test]
fn concurrent_reads_see_a_consistent_prefix_under_heavy_appends() {
let (coord, handle, _dir) = coordinator();
let stop = Arc::new(AtomicBool::new(false));
let readers: Vec<_> = (0..4)
.map(|_| {
let reader = handle.reader();
let stop = Arc::clone(&stop);
thread::spawn(move || {
let mut observed = 0u64;
while !stop.load(Ordering::Relaxed) {
let mut reads = reader.read(&Query::all(), Position::ZERO, None);
let watermark = reads.watermark().get();
let mut positions = Vec::new();
while let Some(item) = reads.next() {
positions.push(item.expect("read failed").position.get());
}
let expected: Vec<u64> = (1..=watermark).collect();
assert_eq!(
positions, expected,
"read did not return the dense prefix 1..={watermark}"
);
observed = observed.max(watermark);
}
observed
})
})
.collect();
let mut rng = Rng(0xABCD_1234_5678_9F01);
for _ in 0..600 {
handle.append(vec![random_event(&mut rng)], None).unwrap();
}
stop.store(true, Ordering::Relaxed);
let mut max_seen = 0u64;
for r in readers {
max_seen = max_seen.max(r.join().unwrap());
}
assert!(max_seen > 0, "readers should have observed some appends");
coord.shutdown();
}
#[test]
fn concurrent_reads_of_the_active_index_see_a_consistent_prefix() {
let (coord, handle, _dir) = coordinator();
let stop = Arc::new(AtomicBool::new(false));
let readers: Vec<_> = (0..4)
.map(|_| {
let reader = handle.reader();
let stop = Arc::clone(&stop);
thread::spawn(move || {
let mut max_seen = 0u64;
while !stop.load(Ordering::Relaxed) {
let query = Query::item(QueryItem::with_tags(tags(&["course:c1"])));
let mut reads = reader.read(&query, Position::ZERO, None);
let watermark = reads.watermark().get();
let mut positions = Vec::new();
while let Some(item) = reads.next() {
positions.push(item.expect("read failed").position.get());
}
let expected: Vec<u64> = (1..=watermark).collect();
assert_eq!(
positions, expected,
"active-index read did not return the dense prefix 1..={watermark}"
);
max_seen = max_seen.max(watermark);
}
max_seen
})
})
.collect();
for _ in 0..600 {
handle
.append(vec![tagged_event("Enrolled", &["course:c1"])], None)
.unwrap();
}
stop.store(true, Ordering::Relaxed);
let mut max_seen = 0u64;
for r in readers {
max_seen = max_seen.max(r.join().unwrap());
}
assert!(max_seen > 0, "readers should have observed some appends");
coord.shutdown();
}
#[test]
fn watermark_resume_has_no_gap_or_duplicate() {
let (coord, handle, _dir) = coordinator();
let mut rng = Rng(0x5EED_5EED_5EED_5EED);
for _ in 0..50 {
handle.append(vec![random_event(&mut rng)], None).unwrap();
}
let first = handle.read(&Query::all(), Position::ZERO, None);
let w1 = first.watermark();
let prefix: Vec<u64> = collect_positions(first);
assert_eq!(prefix, (1..=w1.get()).collect::<Vec<_>>());
for _ in 0..50 {
handle.append(vec![random_event(&mut rng)], None).unwrap();
}
let resumed = handle.read(&Query::all(), w1, None);
let w2 = resumed.watermark();
let tail: Vec<u64> = collect_positions(resumed);
assert!(w2 > w1, "watermark should advance after more appends");
assert_eq!(tail, (w1.get() + 1..=w2.get()).collect::<Vec<_>>());
coord.shutdown();
}
fn collect_positions(mut reads: tephra::read::Reads) -> Vec<u64> {
let mut out = Vec::new();
while let Some(item) = reads.next() {
out.push(item.expect("read failed").position.get());
}
out
}