use super::*;
use crate::block_on;
use crate::protocol::message::encode_frame;
use crate::protocol::transport::poll_fd;
use crate::test_support::{interrupt_self_until, kv_schema, rel, reply_ctrl, reply_status, session_pair, Peer};
use crate::BlockingHost;
use gnitz_wire::control::peek_control_block;
use gnitz_wire::control::ControlHeader;
use gnitz_wire::RelDescriptorBlob;
use gnitz_wire::TypeCode;
use std::collections::HashMap;
use std::os::fd::AsRawFd;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
const TAG: u64 = 0xFEED;
#[derive(Debug, Clone, PartialEq, Eq)]
enum Ev {
Register(u64),
Forget(u64),
Erase(u64),
Fill(u64),
Reseed(u64, u64),
Advance(u64, u64),
}
#[derive(Clone, Default)]
struct Log(Arc<Mutex<Vec<Ev>>>, Arc<Mutex<Option<u64>>>, Arc<Mutex<Option<u64>>>);
impl Log {
fn push(&self, e: Ev) {
self.0.lock().unwrap().push(e);
}
fn take(&self) -> Vec<Ev> {
std::mem::take(&mut *self.0.lock().unwrap())
}
fn fail_next_register_retracting(&self, old: u64) {
*self.1.lock().unwrap() = Some(old);
}
fn fail_next_fill_of(&self, tid: u64) {
*self.2.lock().unwrap() = Some(tid);
}
fn saw(&self, f: impl Fn(&Ev) -> bool) -> bool {
self.0.lock().unwrap().iter().any(f)
}
}
struct StubStore {
cursors: HashMap<u64, DeltaCursor>,
log: Log,
}
impl MirrorStore for StubStore {
fn base_dir(&self) -> &str {
"stub"
}
fn register(&mut self, _name: &RelName, desc: &RelDescriptor) -> Result<(), MirrorError> {
self.log.push(Ev::Register(desc.tid));
if let Some(old) = self.log.1.lock().unwrap().take() {
self.cursors.remove(&old);
return Err(MirrorError::Engine("stub: the copy could not be entered".into()));
}
Ok(())
}
fn forget(&mut self, tid: u64) -> Result<(), MirrorError> {
self.log.push(Ev::Forget(tid));
self.cursors.remove(&tid);
Ok(())
}
fn refill(&mut self, tid: u64) -> Result<(), MirrorError> {
self.log.push(Ev::Erase(tid));
self.cursors.remove(&tid);
Ok(())
}
fn fill(&mut self, tid: u64, blocks: &[&[u8]]) -> Result<(), MirrorError> {
if self.log.2.lock().unwrap().take_if(|failing| *failing == tid).is_some() {
return Err(MirrorError::Engine("stub: the block could not be applied".into()));
}
blocks.iter().for_each(|_| self.log.push(Ev::Fill(tid)));
Ok(())
}
fn seal(&mut self, tid: u64, cursor: DeltaCursor) -> Result<(), MirrorError> {
self.log.push(Ev::Reseed(tid, cursor.tick.get()));
self.cursors.insert(tid, cursor);
Ok(())
}
fn advance(&mut self, tid: u64, _b: &[&[u8]], next: DeltaCursor) -> Result<(), MirrorError> {
if !self.cursors.contains_key(&tid) {
return Err(MirrorError::Engine(format!(
"stub: advance of {tid}, which holds no cursor"
)));
}
self.log.push(Ev::Advance(tid, next.tick.get()));
self.cursors.insert(tid, next);
Ok(())
}
fn scan_spec(
&mut self,
_tid: u64,
_spec: gnitz_wire::ReadSpec,
reply_schema: &Schema,
) -> Result<ZSetBatch, MirrorError> {
Ok(ZSetBatch::new(reply_schema))
}
fn cursor_of(&self, tid: u64) -> Option<DeltaCursor> {
self.cursors.get(&tid).copied()
}
fn checkpoint(&mut self) -> Result<(), MirrorError> {
Ok(())
}
fn poisoned(&self) -> Option<&str> {
None
}
}
const PATIENCE: Duration = Duration::from_secs(5);
impl Peer {
fn expect_frame(&self, what: &str) -> Vec<u8> {
let ready = poll_fd(self.0.as_raw_fd(), libc::POLLIN, Some(Instant::now() + PATIENCE));
assert!(
ready.is_ok_and(|revents| revents & libc::POLLIN != 0),
"{what}: the request never arrived"
);
self.recv()
}
fn expect_request(&self, what: &str) -> u64 {
let frame = self.expect_frame(what);
peek_control_block(&frame).expect("a control header").hdr.target_id
}
fn expect_poll(&self, what: &str) -> Vec<u64> {
self.expect_kept(what).into_iter().map(|(tid, _)| tid).collect()
}
fn expect_kept(&self, what: &str) -> Vec<(u64, u64)> {
let frame = self.expect_frame(what);
let ctrl = peek_control_block(&frame).expect("a control header");
assert_eq!(ctrl.hdr.flags.verb, gnitz_wire::ClientVerb::DeltaPoll, "{what}");
let first = ctrl.hdr.arg0;
let items = gnitz_wire::txn_frame::decode_delta_items(&frame[ctrl.body.clone()]).expect("a delta poll");
(first..).zip(items).map(|(id, item)| (item.view.tid, id)).collect()
}
fn expect_sync_naming(&self, what: &str) -> Vec<u64> {
let frame = self.expect_frame(what);
let ctrl = peek_control_block(&frame).expect("a control header");
assert_eq!(ctrl.hdr.flags.verb, gnitz_wire::ClientVerb::SyncPushed, "{what}");
let mut held = gnitz_wire::txn_frame::decode_held(&frame[ctrl.blob]).expect("held ids");
held.sort_unstable();
held
}
fn reply_watermark(&self, target_id: u64, tag: u64, tick: u64) {
let h = ControlHeader {
target_id,
arg0: tick,
arg1: tag,
..Default::default()
};
self.send(&encode_frame(h, &[], None, None));
}
fn expect_verb(&self, what: &str) -> (gnitz_wire::ClientVerb, ControlHeader) {
let frame = self.expect_frame(what);
let hdr = peek_control_block(&frame).expect("a control header").hdr;
(hdr.flags.verb, hdr)
}
fn expect_sync(&self, what: &str) {
assert_eq!(self.expect_verb(what).0, gnitz_wire::ClientVerb::SyncPushed, "{what}");
}
fn push_train(&self, view: u64, sub: u64, rows: &[(u64, i64, i64)], tick: Option<u64>) {
self.send(&crate::test_support::pushed_marker(view, sub));
let Some(tick) = tick else {
return self.send(&reply_status(view, WireStatus::Error, "lagged"));
};
if !rows.is_empty() {
let hdr = ControlHeader {
target_id: view,
flags: gnitz_wire::WireFlags { continuation: true, ..Default::default() },
..Default::default()
};
self.send(&encode_frame(hdr, &[], None, Some(&crate::test_support::kv_rows(rows))));
}
self.reply_watermark(view, TAG, tick);
}
fn reply_resolved(&self, tid: u64) {
let schema = kv_schema(TypeCode::I64);
let blob = RelDescriptorBlob {
class: RelClass::FedView,
..Default::default()
};
let hdr = ControlHeader { target_id: tid, ..Default::default() };
self.send(&encode_frame(
hdr,
&blob.encode(),
Some(&schema.to_block()),
Some(&ZSetBatch::new(&schema)),
));
}
}
fn fixture(views: &[(u64, &str, u64)]) -> (GnitzClient, Peer, Log) {
let (session, peer) = session_pair();
let mut client = GnitzClient::from_session(session);
let log = Log::default();
let cursors = views
.iter()
.filter_map(|&(tid, _, tick)| Some((tid, DeltaCursor::from_pair(TAG, tick)?)))
.collect();
client.attach_mirror(StubStore { cursors, log: log.clone() }).unwrap();
let schema = Arc::new(kv_schema(TypeCode::I64));
for &(tid, name, _) in views {
let desc = RelDescriptor {
tid,
class: RelClass::FedView,
pk_repeats: false,
serial: false,
schema: Arc::clone(&schema),
indexes: Vec::new(),
token: 0,
};
let entry = MirroredView::new(rel("s", name), Subscription::whole(Arc::new(desc)), None).unwrap();
block_on(client.bind(entry)).unwrap();
client.mirror.as_mut().unwrap().views.get_mut(&tid).unwrap().confirmed = true;
}
assert!(
log.take().iter().all(|e| matches!(e, Ev::Register(_))),
"registration tears nothing down"
);
(client, peer, log)
}
#[test]
fn one_poll_writes_one_request_for_every_view() {
for m in [6usize, 200] {
let names: Vec<String> = (0..m).map(|i| format!("v{i}")).collect();
let views: Vec<(u64, &str, u64)> = names
.iter()
.enumerate()
.map(|(i, n)| (100 + i as u64, n.as_str(), 4))
.collect();
let (mut client, peer, _log) = fixture(&views);
let h = std::thread::spawn(move || {
let ids = peer.expect_poll("the poll");
for &id in &ids {
peer.reply_watermark(id, TAG, 9);
}
ids
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("every view advances");
let seen = h.join().unwrap();
assert_eq!(
client.requests_sent(),
1,
"{m} views: one poll reads every copy and keeps each subscribed"
);
let mut ids = seen;
ids.sort_unstable();
assert_eq!(
ids,
(100..100 + m as u64).collect::<Vec<_>>(),
"every view rides a request"
);
assert_eq!(report.len(), m, "one entry per view");
assert!(
report
.iter()
.all(|o| matches!(o.result, PollResult::Advanced) && o.cursor.is_some_and(|c| c.tick.get() == 9)),
"{m} views: {report:?}"
);
}
}
#[test]
fn only_a_stale_registration_pays_a_resolve() {
for status in [WireStatus::StaleCatalog, WireStatus::DeltaExpired, WireStatus::Error] {
let (mut client, peer, _log) = fixture(&[(7, "v", 4)]);
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_poll("the poll"), vec![7]);
peer.send(&reply_status(7, status, "refused"));
match status {
WireStatus::StaleCatalog => {
assert_eq!(peer.expect_request("the re-resolve"), 0);
peer.send(&reply_ctrl(0, 0));
}
WireStatus::DeltaExpired => {
assert_eq!(peer.expect_poll("the whole read"), vec![7]);
peer.reply_watermark(7, TAG, 20);
}
_ => {}
}
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("a per-view failure is not the call's");
h.join().unwrap();
let (requests, reseeded) = match status {
WireStatus::StaleCatalog => (2, false),
WireStatus::DeltaExpired => (2, true),
_ => (1, false),
};
assert_eq!(client.requests_sent(), requests, "status {status:?}");
let [o] = report.as_slice() else {
panic!("status {status:?}: got {report:?}");
};
assert_eq!(o.view_id, 7);
assert_eq!(o.result.reseeded(), reseeded, "status {status:?}: got {report:?}");
assert_eq!(
matches!(o.result, PollResult::Failed(_)),
!reseeded,
"status {status:?}: got {report:?}"
);
assert!(
client.mirrors(7),
"status {status:?}: a failed poll leaves the copy answering its last round",
);
}
}
#[test]
fn a_bootstrap_fills_the_store_frame_by_frame() {
let (mut client, peer, log) = fixture(&[(7, "a", 0)]);
let watched = log.clone();
let block = |rows: &[(u64, i64, i64)], tick: u64, continuation: bool| {
let hdr = ControlHeader {
target_id: 7,
flags: gnitz_wire::WireFlags { continuation, ..Default::default() },
arg0: tick,
arg1: TAG,
..Default::default()
};
encode_frame(hdr, &[], None, Some(&crate::test_support::kv_rows(rows)))
};
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_poll("the bootstrap"), vec![7]);
peer.send(&block(&[(1, 10, 1)], 0, true));
let until = Instant::now() + PATIENCE;
while !watched.saw(|e| matches!(e, Ev::Fill(7))) {
assert!(Instant::now() < until, "the first block never reached the store");
std::thread::sleep(Duration::from_millis(1));
}
peer.send(&block(&[(2, 20, 1)], 20, false));
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("the bootstrap lands");
h.join().unwrap();
assert!(
matches!(report.as_slice(), [o] if o.view_id == 7 && o.result.reseeded()),
"{report:?}"
);
assert_eq!(
log.take(),
[
Ev::Register(7),
Ev::Erase(7),
Ev::Fill(7),
Ev::Fill(7),
Ev::Reseed(7, 20),
],
);
}
#[test]
fn copies_without_a_cursor_are_read_whole_in_the_one_request() {
let (mut client, peer, log) = fixture(&[(7, "a", 0), (8, "b", 0), (9, "c", 4)]);
let h = std::thread::spawn(move || {
let ids = peer.expect_poll("the poll");
for &id in &ids {
peer.reply_watermark(id, TAG, 20);
}
ids.len()
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("every view is read");
assert_eq!(h.join().unwrap(), 3, "every view rides the request");
assert_eq!(client.requests_sent(), 1);
let mut reseeded: Vec<(u64, bool)> = report.iter().map(|o| (o.view_id, o.result.reseeded())).collect();
reseeded.sort_unstable();
assert_eq!(reseeded, [(7, true), (8, true), (9, false)], "{report:?}");
let mut events = log.take();
events.sort_unstable_by_key(|e| format!("{e:?}"));
assert_eq!(
events,
[
Ev::Advance(9, 20),
Ev::Erase(7),
Ev::Erase(8),
Ev::Register(7),
Ev::Register(8),
Ev::Reseed(7, 20),
Ev::Reseed(8, 20),
],
);
}
#[test]
fn a_refused_block_fails_its_view_and_no_other() {
let (mut client, peer, log) = fixture(&[(7, "a", 0), (8, "b", 4)]);
log.fail_next_fill_of(7);
let block = |rows: &[(u64, i64, i64)]| {
let hdr = ControlHeader {
target_id: 7,
flags: gnitz_wire::WireFlags { continuation: true, ..Default::default() },
..Default::default()
};
encode_frame(hdr, &[], None, Some(&crate::test_support::kv_rows(rows)))
};
let h = std::thread::spawn(move || {
let mut kept = 0;
for (tid, id) in peer.expect_kept("the poll") {
if tid == 7 {
peer.send(&block(&[(1, 10, 1)]));
peer.send(&block(&[(2, 20, 1)]));
} else {
kept = id;
}
peer.reply_watermark(tid, TAG, 20);
}
(peer, kept)
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("a store refusal is the view's failure");
let (peer, kept) = h.join().unwrap();
drop(client.begin_sync(Duration::ZERO).expect("a sync is sent"));
client.session.step(crate::Interest::WRITE);
assert_eq!(peer.expect_sync_naming("the next sync"), [kept], "8's, and not 7's");
let failed: Vec<u64> = report
.iter()
.filter(|o| matches!(o.result, PollResult::Failed(ClientError::Mirror(_))))
.map(|o| o.view_id)
.collect();
assert_eq!(failed, [7], "{report:?}");
assert!(log.saw(|e| *e == Ev::Advance(8, 20)), "the view behind it is read");
assert!(!log.saw(|e| matches!(e, Ev::Fill(7) | Ev::Reseed(7, _))));
}
#[test]
fn a_leftover_poll_does_not_shift_the_replies() {
let (mut client, peer, log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
let from = DeltaCursor::from_pair(TAG, 4);
let abandoned = DeltaPollItem {
view: 7.into(),
from,
reply_layout: kv_schema(TypeCode::I64).layout().layout_digest(),
spec: &[],
};
let first = client.next_sub;
client.next_sub += 1;
drop(DeltaPoll::start(&mut client.session, first, &[abandoned]));
let h = std::thread::spawn(move || {
let abandoned = peer.expect_kept("the abandoned poll");
assert_eq!(abandoned.iter().map(|&(tid, _)| tid).collect::<Vec<_>>(), [7]);
let ids = peer.expect_poll("the poll");
peer.reply_watermark(7, TAG, 50);
for &tid in &ids {
peer.reply_watermark(tid, TAG, 100 + tid);
}
});
block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("the leftover train is drained, not applied");
h.join().unwrap();
let mut applied = log.take();
applied.sort_unstable_by_key(|e| format!("{e:?}"));
assert_eq!(
applied,
[Ev::Advance(7, 107), Ev::Advance(8, 108)],
"each view took the reply its own request opened, and nothing else",
);
}
#[test]
fn no_recovery_runs_before_every_ingest_has() {
let (mut client, peer, log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
let h = std::thread::spawn(move || {
let ids = peer.expect_poll("the poll");
peer.send(&reply_status(ids[0], WireStatus::StaleCatalog, "stale"));
peer.reply_watermark(ids[1], TAG, 11);
assert_eq!(peer.expect_request("the re-resolve"), 0);
peer.reply_resolved(9);
assert_eq!(peer.expect_poll("the next round"), vec![9]);
peer.reply_watermark(9, TAG, 20);
(ids[0], ids[1])
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("a recovered view is not the call's failure");
let (gone, alive) = h.join().unwrap();
assert_eq!(
client.requests_sent(),
3,
"both views ride one request, and the moved one is read whole in another"
);
let events = log.take();
let first_ingest = events
.iter()
.position(|e| matches!(e, Ev::Advance(t, 11) if *t == alive))
.expect("the round ingested the reply it had");
let first_recovery = events.iter().position(|e| !matches!(e, Ev::Advance(..)));
assert!(
first_recovery.is_some_and(|t| t > first_ingest),
"a recovery ran before its round finished ingesting: {events:?}",
);
assert!(
events.contains(&Ev::Reseed(9, 20)) && !events.iter().any(|e| matches!(e, Ev::Erase(t) if *t == gone)),
"the moved view is read whole at its new id: {events:?}",
);
let mut ids: Vec<u64> = report.iter().map(|o| o.view_id).collect();
ids.sort_unstable();
assert_eq!(ids, vec![alive, 9], "one entry per view, the moved one at its new id");
}
#[test]
fn a_reseed_onto_a_live_copy_does_not_erase_it() {
let (mut client, peer, log) = fixture(&[(7, "a", 0), (8, "b", 4)]);
let h = std::thread::spawn(move || {
let mut kept = 0;
for (tid, id) in peer.expect_kept("the poll") {
match tid {
8 => {
kept = id;
peer.reply_watermark(8, TAG, 11)
}
_ => peer.send(&reply_status(7, WireStatus::StaleCatalog, "stale")),
}
}
assert_eq!(peer.expect_request("the re-resolve"), 0);
peer.reply_resolved(8);
let next = peer.expect_kept("the next round");
assert!(matches!(next[..], [(8, id)] if id != kept), "{next:?}");
peer.reply_watermark(8, TAG, 12);
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("the recovery lands on a live copy");
h.join().unwrap();
assert_eq!(client.requests_sent(), 3, "the poll, the resolve, the next round");
let events = log.take();
assert!(
!events.iter().any(|e| matches!(e, Ev::Erase(8) | Ev::Reseed(8, _))),
"the live copy must be neither erased nor re-read whole: {events:?}",
);
assert!(
matches!(report.as_slice(), [o] if o.view_id == 8 && o.cursor.is_some_and(|c| c.tick.get() == 12)),
"one entry, at the id the view ended up under: {report:?}",
);
assert_eq!(client.mirrored_ids(), [8], "the retired registration is gone");
}
#[test]
fn a_recovery_onto_an_expired_view_reads_it_whole_once() {
let (mut client, peer, log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
let h = std::thread::spawn(move || {
for id in peer.expect_poll("the poll") {
let status = if id == 7 {
WireStatus::StaleCatalog
} else {
WireStatus::DeltaExpired
};
peer.send(&reply_status(id, status, "refused"));
}
assert_eq!(peer.expect_request("the re-resolve"), 0);
peer.reply_resolved(8);
assert_eq!(peer.expect_poll("the whole read"), vec![8]);
peer.reply_watermark(8, TAG, 20);
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both recoveries land");
h.join().unwrap();
assert_eq!(client.requests_sent(), 3, "the poll, the resolve, the whole read");
let events = log.take();
let whole: Vec<&Ev> = events
.iter()
.filter(|e| matches!(e, Ev::Erase(_) | Ev::Reseed(..)))
.collect();
assert_eq!(whole, [&Ev::Erase(8), &Ev::Reseed(8, 20)], "{events:?}");
assert!(
matches!(report.as_slice(), [o] if o.view_id == 8 && o.result.reseeded()),
"{report:?}"
);
}
#[test]
fn a_recovery_that_moves_and_then_fails_is_reported_at_the_new_id() {
let (mut client, peer, _log) = fixture(&[(7, "a", 4)]);
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_poll("the poll"), vec![7]);
peer.send(&reply_status(7, WireStatus::StaleCatalog, "stale"));
assert_eq!(peer.expect_request("the re-resolve"), 0);
peer.reply_resolved(9);
assert_eq!(peer.expect_poll("the whole read"), vec![9]);
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("a lost connection is each view's failure");
h.join().unwrap();
assert!(
matches!(
report.as_slice(),
[o] if o.view_id == 9 && o.cursor.is_none() && matches!(o.result, PollResult::Failed(ClientError::ConnectionLost(_)))
),
"{report:?}"
);
assert_eq!(client.mirrored_ids(), [9]);
assert!(!client.mirrors(9), "a copy never synced answers no read");
}
#[test]
fn a_registration_the_store_lost_is_entered_again_by_the_next_poll() {
let (mut client, peer, log) = fixture(&[(7, "a", 4)]);
log.fail_next_register_retracting(7);
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_poll("the poll"), vec![7]);
peer.send(&reply_status(7, WireStatus::StaleCatalog, "stale"));
assert_eq!(peer.expect_request("the re-resolve"), 0);
peer.reply_resolved(9);
peer
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("a store refusal is the view's failure");
let peer = h.join().unwrap();
assert!(
matches!(report.as_slice(), [o] if o.view_id == 7 && matches!(o.result, PollResult::Failed(ClientError::Mirror(_)))),
"{report:?}"
);
assert_eq!(log.take(), [Ev::Register(9)]);
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_poll("the whole read"), vec![7]);
peer.reply_watermark(7, TAG, 20);
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("the view is read whole");
h.join().unwrap();
assert!(
matches!(report.as_slice(), [o] if o.view_id == 7 && o.result.reseeded()),
"{report:?}"
);
assert_eq!(log.take(), [Ev::Register(7), Ev::Erase(7), Ev::Reseed(7, 20)]);
}
#[test]
fn an_interrupted_poll_reannounces_its_reseed() {
let (mut client, peer, log) = fixture(&[(7, "a", 0), (8, "b", 4)]);
client.mirror.as_mut().unwrap().views.get_mut(&7).unwrap().confirmed = false;
let watched = log.clone();
client.host = Box::new(BlockingHost::with_hook(Box::new(move || {
if watched.saw(|e| matches!(e, Ev::Reseed(7, _))) {
Err("interrupted".into())
} else {
Ok(())
}
})));
client.host.attach(client.session.as_fd()).unwrap();
let stop = Arc::new(AtomicBool::new(false));
let sig = interrupt_self_until(Arc::clone(&stop));
let h = std::thread::spawn(move || {
for id in peer.expect_poll("the poll") {
match id {
7 => peer.reply_watermark(7, TAG, 20),
_ => peer.send(&reply_status(8, WireStatus::StaleCatalog, "stale")),
}
}
peer.expect_request("the re-resolve");
peer
});
let r = block_on(client.sync(Duration::ZERO)).map(|s| s.mirrored);
assert!(matches!(r, Err(ClientError::Interrupted(_))), "{r:?}");
stop.store(true, Ordering::Relaxed);
sig.join().unwrap();
client.host = Box::new(BlockingHost::default());
client.host.attach(client.session.as_fd()).unwrap();
let peer = h.join().unwrap();
let h = std::thread::spawn(move || {
peer.send(&reply_ctrl(0, 0)); peer.expect_sync("the first poll");
peer.send(&reply_ctrl(0, 0));
assert_eq!(peer.expect_poll("a poll"), vec![8]);
peer.reply_watermark(8, TAG, 21);
peer.expect_sync("the second poll");
peer.send(&reply_ctrl(0, 0));
});
let mut results = Vec::new();
for _ in 0..2 {
let mut report: Vec<(u64, bool)> = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both views advance")
.iter()
.map(|o| (o.view_id, o.result.reseeded()))
.collect();
report.sort_unstable();
results.push(report);
}
h.join().unwrap();
assert_eq!(
results,
[vec![(7, true), (8, false)], vec![(7, false), (8, false)]],
"the owed reseed is announced once, on the next report that reaches a caller",
);
assert!(
client.cursor_of(7).is_some(),
"7 was read whole and advanced since: its copy answers"
);
assert!(
client.cursor_of(8).is_some(),
"8 was read and subscribed: its copy answers"
);
}
#[test]
fn a_broken_poll_reply_fails_the_views_it_never_answered() {
type Script = fn(&Peer, &[u64]);
let cases: [(&str, Script, usize, bool); 3] = [
("misdirected", |p, ids| p.reply_watermark(ids[1], TAG, 11), 0, true),
("short", |p, ids| p.reply_watermark(ids[0], TAG, 11), 1, true),
(
"a request-level fault",
|p, _| p.send(&reply_status(0, WireStatus::Error, "refused")),
0,
false,
),
];
for (what, script, answered, closes) in cases {
let (mut client, peer, log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
let h = std::thread::spawn(move || {
let ids = peer.expect_poll("the poll");
script(&peer, &ids);
ids
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("a broken reply is reported per view");
let ids = h.join().unwrap();
for (i, id) in ids.iter().enumerate() {
let o = report.iter().find(|o| o.view_id == *id).expect("an entry per view");
if i < answered {
assert!(
matches!(o.result, PollResult::Advanced) && o.cursor.is_some_and(|c| c.tick.get() == 11),
"{what}: the view that was answered keeps its round: {report:?}",
);
} else {
assert!(matches!(o.result, PollResult::Failed(_)), "{what}: {report:?}");
assert!(
!log.saw(|e| matches!(e, Ev::Reseed(t, _) | Ev::Advance(t, _) if t == id)),
"{what}: nothing is ingested under an unanswered view",
);
}
}
assert_eq!(client.session.is_closed(), closes, "{what}");
}
}
#[test]
fn each_view_is_ingested_before_the_next_one_is_answered() {
let (mut client, peer, log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
let watched = log.clone();
let h = std::thread::spawn(move || {
let ids = peer.expect_poll("the poll");
peer.reply_watermark(ids[0], TAG, 11);
let deadline = Instant::now() + PATIENCE;
while !watched.saw(|e| matches!(e, Ev::Reseed(t, _) | Ev::Advance(t, _) if *t == ids[0])) {
assert!(
Instant::now() < deadline,
"the first view's blocks must reach the store before the second view is answered",
);
std::thread::sleep(Duration::from_millis(1));
}
peer.reply_watermark(ids[1], TAG, 12);
});
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both views advance");
h.join().unwrap();
assert!(report.iter().all(|o| matches!(o.result, PollResult::Advanced)));
}
#[test]
fn a_dead_connection_fails_each_view_rather_than_the_call() {
for lost_before in [false, true] {
let (mut client, peer, _log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
drop(peer);
if lost_before {
let spec = gnitz_wire::ReadSpec::all_rows(gnitz_wire::ReadBound::None);
let r = block_on(client.scan_spec(1, &spec, &Arc::new(kv_schema(TypeCode::I64))));
assert!(matches!(r, Err(ClientError::ConnectionLost(_))), "{r:?}");
}
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("a dead socket is not the call's own failure");
let mut ids: Vec<u64> = report.iter().map(|o| o.view_id).collect();
ids.sort_unstable();
assert_eq!(ids, vec![7, 8], "one entry per view: {report:?}");
assert!(
report
.iter()
.all(|o| matches!(o.result, PollResult::Failed(ClientError::ConnectionLost(_)))),
"lost before: {lost_before}: {report:?}",
);
assert!(
client.mirrors(7) && client.mirrors(8),
"and both copies still answer at the round they reached",
);
}
}
fn subscribed() -> (GnitzClient, Peer, Log, HashMap<u64, u64>) {
let (mut client, peer, log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
let h = std::thread::spawn(move || {
let kept = peer.expect_kept("the poll");
for &(id, _) in &kept {
peer.reply_watermark(id, TAG, 9);
}
(peer, kept.into_iter().collect())
});
block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both views advance");
let (peer, subs) = h.join().unwrap();
log.take();
(client, peer, log, subs)
}
#[test]
fn a_subscribed_copy_is_advanced_by_what_was_pushed() {
let (mut client, peer, log, subs) = subscribed();
let h = std::thread::spawn(move || {
peer.expect_sync("the poll");
peer.push_train(7, subs[&7], &[(1, 10, 1)], Some(11));
peer.push_train(7, subs[&7], &[(1, 10, -1)], Some(12));
peer.push_train(8, subs[&8], &[], Some(12));
peer.send(&reply_ctrl(0, 0));
});
let before = client.requests_sent();
let mut report: Vec<(u64, u64)> = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both views advance")
.iter()
.map(|o| {
assert!(matches!(o.result, PollResult::Advanced), "{o:?}");
(o.view_id, o.cursor.expect("a copy that answers").tick.get())
})
.collect();
h.join().unwrap();
report.sort_unstable();
assert_eq!(client.requests_sent() - before, 1, "one request for both views");
assert_eq!(report, [(7, 12), (8, 12)]);
let mut applied = log.take();
applied.sort_unstable_by_key(|e| format!("{e:?}"));
assert_eq!(applied, [Ev::Advance(7, 12), Ev::Advance(8, 12)]);
}
#[test]
fn a_train_behind_another_reply_waits_for_its_poll() {
let (mut client, peer, log, subs) = subscribed();
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_request("an id allocation"), 0);
peer.send(&reply_ctrl(0, 4242));
peer.push_train(8, subs[&8], &[(2, 20, 1)], Some(15));
peer.expect_sync("the poll");
peer.send(&reply_ctrl(0, 0));
});
assert_eq!(block_on(client.alloc_id()).unwrap(), 4242);
assert!(log.take().is_empty(), "a copy moves only in a poll");
block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both views advance");
h.join().unwrap();
assert_eq!(log.take(), [Ev::Advance(8, 15)], "a view sent nothing does not move");
}
#[test]
fn a_subscription_that_ended_is_continued_by_a_delta_read() {
let (mut client, peer, log, subs) = subscribed();
let h = std::thread::spawn(move || {
peer.expect_sync("the poll");
peer.push_train(7, subs[&7], &[], None);
peer.send(&reply_ctrl(0, 0));
let again = peer.expect_kept("the delta read");
assert!(matches!(again[..], [(7, id)] if id != 0 && id != subs[&7]), "{again:?}");
peer.reply_watermark(7, TAG, 13);
});
let before = client.requests_sent();
let report = block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both views advance");
h.join().unwrap();
assert_eq!(client.requests_sent() - before, 2, "the sync, then the delta read");
assert!(
report.iter().all(|o| matches!(o.result, PollResult::Advanced)),
"an ended subscription is nobody's failure: {report:?}"
);
assert_eq!(log.take(), [Ev::Advance(7, 13)]);
}
#[test]
fn one_sync_serves_a_reader_and_the_mirror() {
let (mut client, peer, log, subs) = subscribed();
let whole = gnitz_wire::ReadSpec::all_rows(gnitz_wire::ReadBound::None).encode();
let from = DeltaCursor::from_pair(TAG, 9).unwrap();
let sub = client
.subscribe(Target::from(40), from, &Arc::new(kv_schema(TypeCode::I64)), &whole)
.expect("the request is sent");
let h = std::thread::spawn(move || {
let reader = peer.expect_kept("the reader's subscription");
assert_eq!(reader, [(40, sub)]);
peer.reply_watermark(40, TAG, 9);
peer.expect_sync("the sync");
peer.push_train(7, subs[&7], &[(1, 10, 1)], Some(11));
peer.push_train(40, sub, &[(2, 20, 1)], Some(11));
peer.send(&reply_ctrl(0, 0));
});
let before = client.requests_sent();
let synced = block_on(client.sync(Duration::ZERO)).expect("one sync");
h.join().unwrap();
assert_eq!(client.requests_sent() - before, 1, "one request for both");
let [pushed] = &synced.pushed[..] else {
panic!("one entry per subscription")
};
let (rows, at) = pushed.result.as_ref().expect("the subscription is held");
assert_eq!((pushed.sub, rows.batch.len(), at.tick.get()), (sub, 1, 11));
assert_eq!(synced.mirrored.len(), 2, "one entry per view");
assert_eq!(log.take(), [Ev::Advance(7, 11)]);
let again = client.session.take_pushed(sub).expect("the subscription is held");
assert!(again.0.is_empty(), "one sync takes everything it brought");
}
#[test]
fn an_interrupted_sync_keeps_a_readers_deltas() {
let (mut client, peer, _log) = fixture(&[(7, "a", 4)]);
let whole = gnitz_wire::ReadSpec::all_rows(gnitz_wire::ReadBound::None).encode();
let from = DeltaCursor::from_pair(TAG, 9).unwrap();
let sub = client
.subscribe(Target::from(40), from, &Arc::new(kv_schema(TypeCode::I64)), &whole)
.expect("the request is sent");
let reading = Arc::new(AtomicBool::new(false));
let watched = Arc::clone(&reading);
client.host = Box::new(BlockingHost::with_hook(Box::new(move || {
if watched.load(Ordering::Relaxed) {
Err("interrupted".into())
} else {
Ok(())
}
})));
client.host.attach(client.session.as_fd()).unwrap();
let stop = Arc::new(AtomicBool::new(false));
let sig = interrupt_self_until(Arc::clone(&stop));
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_kept("the reader's subscription"), [(40, sub)]);
peer.reply_watermark(40, TAG, 9);
peer.expect_sync("the sync");
peer.push_train(40, sub, &[(2, 20, 1)], Some(11));
peer.send(&reply_ctrl(0, 0));
let read = peer.expect_kept("the mirror's read");
assert!(matches!(read[..], [(7, id)] if id != 0), "{read:?}");
reading.store(true, Ordering::Relaxed);
(peer, read[0].1)
});
let r = block_on(client.sync(Duration::ZERO)).map(|s| s.mirrored);
assert!(matches!(r, Err(ClientError::Interrupted(_))), "{r:?}");
stop.store(true, Ordering::Relaxed);
sig.join().unwrap();
client.host = Box::new(BlockingHost::default());
client.host.attach(client.session.as_fd()).unwrap();
let (peer, abandoned) = h.join().unwrap();
let h = std::thread::spawn(move || {
peer.reply_watermark(7, TAG, 11);
let held = peer.expect_sync_naming("the next sync");
assert!(held == [sub] && sub != abandoned, "{held:?}");
peer.send(&reply_ctrl(0, 0));
assert_eq!(peer.expect_poll("the mirror's read"), vec![7]);
peer.reply_watermark(7, TAG, 11);
});
let synced = block_on(client.sync(Duration::ZERO)).expect("the next sync");
h.join().unwrap();
let [pushed] = &synced.pushed[..] else {
panic!("one entry per subscription")
};
let (rows, at) = pushed.result.as_ref().expect("the subscription is held");
assert_eq!((pushed.sub, rows.batch.len(), at.tick.get()), (sub, 1, 11));
}
#[test]
fn a_view_that_keeps_failing_does_not_end_the_hold() {
let (mut client, peer, _log) = fixture(&[(7, "a", 4), (8, "b", 4)]);
let h = std::thread::spawn(move || {
for id in peer.expect_poll("the first read") {
match id {
7 => peer.reply_watermark(7, TAG, 9),
_ => peer.send(&reply_status(id, WireStatus::Error, "gone")),
}
}
let (verb, hdr) = peer.expect_verb("the second sync");
assert_eq!(verb, gnitz_wire::ClientVerb::SyncPushed);
peer.send(&reply_ctrl(0, 0));
assert_eq!(peer.expect_poll("the retry"), vec![8]);
peer.send(&reply_status(8, WireStatus::Error, "gone"));
hdr.arg0
});
let failed = |synced: crate::Synced| {
let failed = synced
.mirrored
.iter()
.filter(|o| matches!(o.result, PollResult::Failed(_)));
failed.count()
};
assert_eq!(failed(block_on(client.sync(Duration::ZERO)).unwrap()), 1);
assert_eq!(failed(block_on(client.sync(Duration::from_secs(30))).unwrap()), 1);
assert_eq!(h.join().unwrap(), 30_000, "the sync asked to be held");
}
#[test]
fn a_sync_over_an_untaken_train_asks_for_no_hold() {
let (mut client, peer, log, subs) = subscribed();
let h = std::thread::spawn(move || {
assert_eq!(peer.expect_request("an id allocation"), 0);
peer.push_train(8, subs[&8], &[(2, 20, 1)], Some(15));
peer.send(&reply_ctrl(0, 4242));
let (_, hdr) = peer.expect_verb("the sync");
peer.send(&reply_ctrl(0, 0));
hdr.arg0
});
assert_eq!(block_on(client.alloc_id()).unwrap(), 4242);
let synced = block_on(client.sync(Duration::from_secs(30))).unwrap();
assert_eq!(synced.mirrored.len(), 2);
assert_eq!(h.join().unwrap(), 0, "no hold in front of a train already here");
assert_eq!(log.take(), [Ev::Advance(8, 15)]);
}
#[test]
fn a_forgotten_view_ends_its_subscription() {
let (mut client, peer, _log, subs) = subscribed();
let h = std::thread::spawn(move || {
let mut both: Vec<u64> = subs.values().copied().collect();
both.sort_unstable();
assert_eq!(peer.expect_sync_naming("the poll"), both);
peer.send(&reply_ctrl(0, 0));
assert_eq!(peer.expect_sync_naming("the poll after"), [subs[&8]]);
peer.send(&reply_ctrl(0, 0));
});
block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("both views advance");
block_on(client.forget_view(7)).unwrap();
block_on(client.sync(Duration::ZERO))
.map(|s| s.mirrored)
.expect("the view kept advances");
h.join().unwrap();
}