use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, mpsc};
use std::task::{Wake, Waker};
use ridl_loopback::{Handles, Loopback};
use ridl_rt::contract::{CatalogHash, CatalogRef, InterfaceNo, Ordinal};
use ridl_rt::error::{CallError, Contract, Transport};
use ridl_rt::port::{
Attached, Caller, ClaimId, Clock, EventSink, EventSource, FixedReader, Handler, Interest,
ReadError, SendError, SettleError, SignalReader, SignalWriter, Wakeable,
};
use ridl_rt::sample::{Duration, Freshness, Provenance, Timestamp};
const IFACE: InterfaceNo = InterfaceNo(1);
const ORD: Ordinal = Ordinal(1);
const OTHER: Ordinal = Ordinal(2);
fn catalog() -> CatalogRef {
CatalogRef {
name: "face.demo",
hash: CatalogHash([0u8; 32]),
}
}
fn runtime() -> Loopback {
Loopback::new(catalog())
}
#[test]
fn two_runtimes_start_at_the_same_logical_time() {
let first = runtime();
std::thread::sleep(std::time::Duration::from_millis(5));
let second = runtime();
assert_eq!(first.now(), Timestamp(0), "the clock starts at 0");
assert_eq!(
first.now(),
second.now(),
"the clock must not read wall-clock time"
);
}
#[test]
#[should_panic(expected = "the clock advances forward")]
fn the_clock_refuses_to_run_backwards() {
let mut rt = runtime();
rt.advance(Duration(-1));
}
#[test]
fn every_value_is_unbounded_because_the_runtime_has_no_member_table() {
let mut rt = runtime();
rt.set(IFACE, ORD, &[1]).expect("set");
rt.commit();
let mut out = [0u8; 8];
assert_eq!(
rt.read(IFACE, ORD, &mut out).expect("read").freshness,
Freshness::Unbounded,
"freshness is measured against a staleness bound this runtime cannot read"
);
}
#[test]
fn a_sink_counts_its_own_channel_and_not_another_sinks() {
let rt = runtime();
let mut source = rt.source();
let mut first = rt.sink();
let mut second = rt.sink();
source.subscribe(IFACE, &[ORD]).expect("subscribe");
first.raise(IFACE, ORD, &[1], None).expect("raise");
second.raise(IFACE, ORD, &[2], None).expect("raise");
first.raise(IFACE, ORD, &[3], None).expect("raise");
let mut out = [0u8; 8];
let mut seqs = Vec::new();
while let Some(occurrence) = source.next(&mut out).expect("next") {
seqs.push(occurrence.envelope.seq);
}
assert_eq!(seqs, vec![1, 1, 2]);
}
#[test]
fn serve_records_what_it_was_asked_to_present() {
let rt = runtime();
let mut handler = rt.handler();
handler.serve(IFACE, &[ORD, OTHER]).expect("serve");
handler.serve(IFACE, &[ORD]).expect("serve again");
assert_eq!(handler.served(), &[(IFACE, ORD), (IFACE, OTHER)]);
}
#[test]
fn a_handler_that_served_nothing_is_presented_every_call() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
caller
.command(InterfaceNo(2), OTHER, &[1], None)
.expect("send");
let mut buf = [0u8; 8];
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(claim.iface, InterfaceNo(2));
}
#[test]
fn the_injected_settle_failure_is_too_large_with_no_capacity() {
let mut rt = runtime();
rt.command(IFACE, ORD, &[1], None).expect("send");
let mut buf = [0u8; 8];
let claim = rt.next_claim(&mut buf).expect("read").expect("waiting");
rt.fail_next_settle();
assert_eq!(
rt.settle(claim.id, Ok(&[])),
Err(SettleError::TooLarge { cap: 0 }),
"the one fault this runtime injects"
);
}
#[test]
fn a_provisioned_fixed_reads_back_and_an_unprovisioned_one_reports_it() {
let mut rt = runtime();
let mut out = [0u8; 8];
assert_eq!(
rt.read_fixed(IFACE, ORD, &mut out),
Err(ReadError::Contract(Contract::UnknownInteraction)),
"nothing was provisioned at that ordinal"
);
rt.provision_fixed(IFACE, ORD, &[1, 2]);
assert_eq!(rt.read_fixed(IFACE, ORD, &mut out).expect("read"), 2);
assert_eq!(&out[..2], &[1, 2]);
let mut short = [0u8; 1];
assert_eq!(
rt.read_fixed(IFACE, ORD, &mut short),
Err(ReadError::Short { needed: 2 })
);
}
#[test]
fn every_split_handle_carries_the_catalog() {
let handles = runtime().split();
assert_eq!(*handles.reader.catalog(), catalog());
assert_eq!(*handles.writer.catalog(), catalog());
assert_eq!(*handles.source.catalog(), catalog());
assert_eq!(*handles.sink.catalog(), catalog());
assert_eq!(*handles.caller.catalog(), catalog());
assert_eq!(*handles.handler.catalog(), catalog());
}
#[test]
fn an_attached_aggregate_carries_the_catalog() {
let attached = runtime().attach();
assert_eq!(*attached.caller().catalog(), catalog());
let handles = attached.split();
assert_eq!(*handles.reader.catalog(), catalog());
assert_eq!(*handles.writer.catalog(), catalog());
assert_eq!(*handles.source.catalog(), catalog());
assert_eq!(*handles.sink.catalog(), catalog());
assert_eq!(*handles.caller.catalog(), catalog());
assert_eq!(*handles.handler.catalog(), catalog());
}
#[test]
fn a_value_committed_through_one_aggregate_is_read_through_an_attached_one() {
let mut rt = runtime();
rt.set(IFACE, ORD, &[1]).expect("staged");
let mut attached = rt.attach();
attached.set(IFACE, OTHER, &[2]).expect("staged");
attached.commit();
let mut out = [0u8; 8];
let unpublished = rt.read(IFACE, ORD, &mut out).expect("read");
assert_eq!(unpublished.provenance, Provenance::Init);
let published = rt.read(IFACE, OTHER, &mut out).expect("read");
assert_eq!(&out[..published.len], &[2]);
rt.commit();
let sample = attached.read(IFACE, ORD, &mut out).expect("read");
assert_eq!(&out[..sample.len], &[1]);
}
#[test]
fn an_attached_aggregate_advances_the_one_clock() {
let rt = runtime();
let mut attached = rt.attach();
attached.advance(Duration(5));
assert_eq!(rt.now(), Timestamp(5));
}
#[test]
fn an_attached_aggregate_starts_with_none_of_the_originals_handle_state() {
let mut rt = runtime();
rt.subscribe(IFACE, &[ORD]).expect("subscribe");
rt.serve(IFACE, &[ORD]).expect("serve");
rt.set(IFACE, ORD, &[1]).expect("staged");
rt.commit();
rt.raise(IFACE, ORD, &[1], None).expect("raise");
rt.command(IFACE, ORD, &[1], None).expect("send");
let mut attached = rt.attach();
let mut out = [0u8; 8];
attached.set(IFACE, ORD, &[2]).expect("staged");
attached.commit();
let sample = attached.read(IFACE, ORD, &mut out).expect("read");
assert_eq!(
sample.envelope.seq, 1,
"the attached writer's first publication"
);
attached.raise(IFACE, ORD, &[2], None).expect("raise");
assert!(
attached.next(&mut out).expect("next").is_none(),
"the attached source is subscribed to nothing"
);
let mut raised = Vec::new();
while let Some(occurrence) = rt.next(&mut out).expect("next") {
raised.push(occurrence.envelope.seq);
}
assert_eq!(
raised,
vec![1, 1],
"each sink's first raise is its own seq 1"
);
attached.command(IFACE, ORD, &[2], None).expect("send");
let mut sent = Vec::new();
while let Some(claim) = rt.next_claim(&mut out).expect("next_claim") {
sent.push(claim.envelope.seq);
}
assert_eq!(
sent,
vec![1, 1],
"each caller's first call is its own seq 1"
);
assert!(attached.split().handler.served().is_empty());
}
#[test]
fn an_attached_aggregate_has_its_own_subscriptions_and_queue() {
let mut rt = runtime();
let mut attached = rt.attach();
rt.subscribe(IFACE, &[ORD]).expect("subscribe");
attached.raise(IFACE, ORD, &[1], None).expect("raise");
let mut out = [0u8; 8];
assert!(
attached.next(&mut out).expect("next").is_none(),
"the attached aggregate's source is subscribed to nothing"
);
attached.subscribe(IFACE, &[ORD]).expect("subscribe");
attached.raise(IFACE, ORD, &[2], None).expect("raise");
let mut drain = |port: &mut Loopback| {
let mut payloads = Vec::new();
while let Some(occurrence) = port.next(&mut out).expect("next") {
payloads.push(out[..occurrence.len].to_vec());
}
payloads
};
assert_eq!(drain(&mut rt), vec![vec![1], vec![2]]);
assert_eq!(drain(&mut attached), vec![vec![2]]);
}
#[test]
fn a_call_sent_through_one_aggregate_is_served_through_an_attached_one() {
let mut rt = runtime();
let mut attached = rt.attach();
let sent = rt.command(IFACE, ORD, &[1], None).expect("send");
let mut buf = [0u8; 8];
let claim = attached
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
attached.settle(claim.id, Ok(&[])).expect("settle");
assert_eq!(rt.ack(sent), Some(Ok(())));
}
#[test]
fn dropping_an_attached_aggregate_leaves_the_originals_calls_and_subscriptions() {
let mut rt = runtime();
rt.subscribe(IFACE, &[ORD]).expect("subscribe");
let sent = rt.command(IFACE, ORD, &[1], None).expect("send");
let mut attached = rt.attach();
attached.subscribe(IFACE, &[ORD]).expect("subscribe");
attached.command(IFACE, ORD, &[2], None).expect("send");
let mut out = [0u8; 8];
let claim = rt
.next_claim(&mut out)
.expect("next_claim")
.expect("the original's call, sent first");
assert_eq!(out[0], 1);
drop(attached);
rt.raise(IFACE, ORD, &[3], None).expect("raise");
let occurrence = rt.next(&mut out).expect("next").expect("subscribed");
assert_eq!(&out[..occurrence.len], &[3]);
rt.settle(claim.id, Ok(&[]))
.expect("the claim is still the original's");
assert_eq!(rt.ack(sent), Some(Ok(())));
assert!(
rt.next_claim(&mut out).expect("next_claim").is_none(),
"the dropped aggregate's call was withdrawn"
);
}
#[test]
fn an_attached_aggregate_keeps_the_store_after_the_original_is_dropped() {
let mut rt = runtime();
let mut attached = rt.attach();
attached.subscribe(IFACE, &[ORD]).expect("subscribe");
let sent = attached.command(IFACE, ORD, &[5], None).expect("send");
rt.set(IFACE, ORD, &[4]).expect("staged");
rt.commit();
drop(rt);
let mut out = [0u8; 8];
let sample = attached.read(IFACE, ORD, &mut out).expect("read");
assert_eq!(&out[..sample.len], &[4]);
attached.raise(IFACE, ORD, &[6], None).expect("raise");
let occurrence = attached.next(&mut out).expect("next").expect("subscribed");
assert_eq!(&out[..occurrence.len], &[6]);
let claim = attached
.next_claim(&mut out)
.expect("next_claim")
.expect("the attached aggregate's call");
assert_eq!(out[0], 5);
attached.settle(claim.id, Ok(&[])).expect("settle");
assert_eq!(attached.ack(sent), Some(Ok(())));
}
#[test]
fn a_writer_handle_publishes_on_one_thread_while_a_reader_reads_on_another() {
let handles = runtime().split();
let reader = handles.reader;
let mut writer = handles.writer;
let publisher = std::thread::spawn(move || {
for value in 1..=50u8 {
writer.set(IFACE, ORD, &[value; 4]).expect("set");
writer.commit();
}
});
let mut out = [0u8; 8];
let mut seen = 0u32;
while seen < 200 {
let raw = reader.read(IFACE, ORD, &mut out).expect("read");
if raw.len == 0 {
continue;
}
assert_eq!(raw.len, 4);
let value = out[0];
assert!((1..=50).contains(&value));
assert_eq!(
&out[..4],
&[value; 4],
"a publication is read whole or not at all"
);
assert_eq!(
raw.envelope.seq,
u64::from(value),
"and its envelope belongs to the value read"
);
seen += 1;
}
publisher.join().expect("the publishing thread finished");
let raw = reader.read(IFACE, ORD, &mut out).expect("read");
assert_eq!(&out[..raw.len], &[50; 4]);
assert_eq!(raw.envelope.seq, 50);
}
#[test]
fn a_reader_handle_is_shared_between_threads() {
let handles = runtime().split();
let mut writer = handles.writer;
writer.set(IFACE, ORD, &[3]).expect("set");
writer.commit();
let reader = std::sync::Arc::new(handles.reader);
let readers: Vec<_> = (0..4)
.map(|_| {
let reader = std::sync::Arc::clone(&reader);
std::thread::spawn(move || {
let mut out = [0u8; 8];
let raw = reader.read(IFACE, ORD, &mut out).expect("read");
out[..raw.len].to_vec()
})
})
.collect();
for reader in readers {
assert_eq!(reader.join().expect("the reading thread finished"), vec![3]);
}
}
struct Count(AtomicUsize);
impl Wake for Count {
fn wake(self: Arc<Self>) {
self.0.fetch_add(1, Ordering::SeqCst);
}
}
fn counting() -> (Arc<Count>, Waker) {
let count = Arc::new(Count(AtomicUsize::new(0)));
(Arc::clone(&count), Waker::from(count))
}
fn wakes(count: &Count) -> usize {
count.0.load(Ordering::SeqCst)
}
#[test]
fn each_settlement_wakes_only_its_own_calls_waiter() {
let Handles {
mut caller,
mut handler,
..
} = runtime().split();
let mut buf = [0u8; 8];
let first_call = caller.command(IFACE, ORD, &[1], None).expect("send");
let second_call = caller.command(IFACE, ORD, &[2], None).expect("send");
let (first, first_waker) = counting();
let (second, second_waker) = counting();
caller.wake_on(Interest::Outcome(first_call), &first_waker);
caller.wake_on(Interest::Outcome(second_call), &second_waker);
assert_eq!(
(wakes(&first), wakes(&second)),
(0, 0),
"two calls, two wakers, no displacement"
);
let claim_one = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
let claim_two = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
handler.settle(claim_two.id, Ok(&[])).expect("settle");
assert_eq!(
(wakes(&first), wakes(&second)),
(0, 1),
"only the second call is settled"
);
handler.settle(claim_one.id, Ok(&[])).expect("settle");
assert_eq!((wakes(&first), wakes(&second)), (1, 1));
}
#[test]
fn a_waiter_registered_after_the_settlement_is_woken_at_once() {
let Handles {
mut caller,
mut handler,
..
} = runtime().split();
let c = caller.query(IFACE, ORD, &[1], None).expect("send");
let mut buf = [0u8; 8];
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
handler.settle(claim.id, Ok(&[7])).expect("settle");
let (count, waker) = counting();
caller.wake_on(Interest::Outcome(c), &waker);
assert_eq!(wakes(&count), 1, "the outcome is already known");
assert_eq!(caller.reply(c, &mut buf), Ok(Some(Ok(1))));
}
#[test]
fn a_raise_wakes_a_subscribed_source_and_not_an_unsubscribed_one() {
let rt = runtime();
let mut sink = rt.sink();
let mut subscribed = rt.source();
let mut elsewhere = rt.source();
subscribed.subscribe(IFACE, &[ORD]).expect("subscribe");
elsewhere.subscribe(IFACE, &[OTHER]).expect("subscribe");
let unsubscribed = rt.source();
let (woken, woken_waker) = counting();
let (other, other_waker) = counting();
let (none, none_waker) = counting();
subscribed.wake_on(Interest::Event(IFACE), &woken_waker);
elsewhere.wake_on(Interest::Event(IFACE), &other_waker);
unsubscribed.wake_on(Interest::Event(IFACE), &none_waker);
sink.raise(IFACE, ORD, &[1], None).expect("raise");
assert_eq!(wakes(&woken), 1, "the subscribed source is woken");
assert_eq!(
wakes(&other),
0,
"a source subscribed to another event is not"
);
assert_eq!(wakes(&none), 0, "an unsubscribed source is not");
let mut out = [0u8; 8];
assert!(subscribed.next(&mut out).expect("next").is_some());
}
#[test]
fn a_waiting_occurrence_of_another_interface_wakes_an_event_registration_at_once() {
let rt = runtime();
let mut sink = rt.sink();
let mut source = rt.source();
source.subscribe(InterfaceNo(2), &[ORD]).expect("subscribe");
sink.raise(InterfaceNo(2), ORD, &[1], None).expect("raise");
let (count, waker) = counting();
source.wake_on(Interest::Event(IFACE), &waker);
assert_eq!(
wakes(&count),
1,
"an occurrence of another interface is waiting"
);
}
#[test]
fn a_waiting_call_on_another_interface_wakes_a_claim_registration_at_once() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
handler.serve(InterfaceNo(2), &[ORD]).expect("serve");
caller
.command(InterfaceNo(2), ORD, &[1], None)
.expect("send");
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
assert_eq!(wakes(&count), 1, "a call on another interface is waiting");
}
#[test]
fn a_drop_returns_only_the_dropped_handlers_claims_and_wakes_every_serving_handler() {
let rt = runtime();
let mut caller = rt.caller();
let mut keeper = rt.handler();
let mut dropped = rt.handler();
let mut other = rt.handler();
let mut elsewhere = rt.handler();
keeper.serve(IFACE, &[ORD]).expect("serve");
dropped.serve(IFACE, &[ORD]).expect("serve");
other.serve(IFACE, &[ORD]).expect("serve");
elsewhere.serve(IFACE, &[OTHER]).expect("serve");
let mut buf = [0u8; 8];
let kept = caller.command(IFACE, ORD, &[1], None).expect("send");
let kept_claim = keeper
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
caller.command(IFACE, ORD, &[2], None).expect("send");
dropped
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
let (first, first_waker) = counting();
let (second, second_waker) = counting();
let (third, third_waker) = counting();
keeper.wake_on(Interest::Claim(IFACE), &first_waker);
other.wake_on(Interest::Claim(IFACE), &second_waker);
elsewhere.wake_on(Interest::Claim(IFACE), &third_waker);
drop(dropped);
assert_eq!(
wakes(&first),
1,
"every serving handler is woken: the first"
);
assert_eq!(
wakes(&second),
1,
"every serving handler is woken: the second"
);
assert_eq!(
wakes(&third),
0,
"a handler serving another member is not woken"
);
assert_eq!(
elsewhere.next_claim(&mut buf).expect("next_claim"),
None,
"and is presented nothing"
);
let returned = other
.next_claim(&mut buf)
.expect("next_claim")
.expect("returned");
assert_eq!(
&buf[..returned.len],
&[2],
"only the dropped handler's call"
);
assert_eq!(
other.next_claim(&mut buf).expect("next_claim"),
None,
"the call another handler holds stays with it"
);
keeper
.settle(kept_claim.id, Ok(&[]))
.expect("the keeper still holds its claim");
assert_eq!(caller.ack(kept), Some(Ok(())));
}
#[test]
fn a_waiter_woken_at_once_is_not_stored() {
let rt = runtime();
let mut sink = rt.sink();
let mut source = rt.source();
source.subscribe(IFACE, &[ORD]).expect("subscribe");
sink.raise(IFACE, ORD, &[1], None).expect("raise");
let (count, waker) = counting();
source.wake_on(Interest::Event(IFACE), &waker);
assert_eq!(wakes(&count), 1, "an occurrence is waiting");
sink.raise(IFACE, ORD, &[2], None).expect("raise");
assert_eq!(wakes(&count), 1, "a waker woken at once was not stored");
}
#[test]
fn a_registration_on_a_forgotten_call_is_woken_at_once() {
let Handles {
mut caller,
mut handler,
..
} = runtime().split();
let c = caller.command(IFACE, ORD, &[1], None).expect("send");
handler
.next_claim(&mut [0u8; 8])
.expect("next_claim")
.expect("waiting");
caller.forget(c);
let (count, waker) = counting();
caller.wake_on(Interest::Outcome(c), &waker);
assert_eq!(wakes(&count), 1, "no outcome will be recorded");
}
#[test]
fn every_subscribed_source_and_every_serving_handler_is_woken() {
let rt = runtime();
let mut sink = rt.sink();
let mut caller = rt.caller();
let mut first_source = rt.source();
let mut second_source = rt.source();
first_source.subscribe(IFACE, &[ORD]).expect("subscribe");
second_source.subscribe(IFACE, &[ORD]).expect("subscribe");
let mut first_handler = rt.handler();
let mut second_handler = rt.handler();
first_handler.serve(IFACE, &[ORD]).expect("serve");
second_handler.serve(IFACE, &[ORD]).expect("serve");
let wakers: Vec<_> = (0..4).map(|_| counting()).collect();
first_source.wake_on(Interest::Event(IFACE), &wakers[0].1);
second_source.wake_on(Interest::Event(IFACE), &wakers[1].1);
first_handler.wake_on(Interest::Claim(IFACE), &wakers[2].1);
second_handler.wake_on(Interest::Claim(IFACE), &wakers[3].1);
sink.raise(IFACE, ORD, &[1], None).expect("raise");
assert_eq!(wakes(&wakers[0].0), 1, "the first source");
assert_eq!(wakes(&wakers[1].0), 1, "the second source");
caller.command(IFACE, ORD, &[1], None).expect("send");
assert_eq!(wakes(&wakers[2].0), 1, "the first handler");
assert_eq!(wakes(&wakers[3].0), 1, "the second handler");
}
#[test]
fn a_claim_registration_is_not_woken_by_a_call_the_handler_does_not_serve() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
handler.serve(IFACE, &[OTHER]).expect("serve");
caller.command(IFACE, ORD, &[1], None).expect("send");
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
assert_eq!(wakes(&count), 0, "the waiting call is not this handler's");
}
#[test]
fn a_serve_that_admits_nothing_does_not_wake_the_handler() {
let rt = runtime();
let mut handler = rt.handler();
handler.serve(IFACE, &[OTHER]).expect("serve");
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
handler.serve(IFACE, &[ORD]).expect("serve");
assert_eq!(wakes(&count), 0, "no call is waiting");
}
#[test]
fn a_dropped_handler_returns_its_claims_to_the_waiting_calls() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
first.serve(IFACE, &[ORD]).expect("serve");
second.serve(IFACE, &[ORD]).expect("serve");
let mut buf = [0u8; 8];
let c = caller.command(IFACE, ORD, &[1], None).expect("send");
first
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
let later = caller.command(IFACE, ORD, &[2], None).expect("send");
first
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
second.next_claim(&mut buf).expect("next_claim"),
None,
"both calls are held by the first handler"
);
let (count, waker) = counting();
second.wake_on(Interest::Claim(IFACE), &waker);
assert_eq!(wakes(&count), 0);
let (outcome, outcome_waker) = counting();
caller.wake_on(Interest::Outcome(c), &outcome_waker);
drop(first);
assert_eq!(
wakes(&count),
1,
"the returned claims wake a serving handler"
);
assert_eq!(wakes(&outcome), 0, "a returned claim is not an outcome");
let again = second
.next_claim(&mut buf)
.expect("next_claim")
.expect("returned");
assert_eq!(&buf[..again.len], &[1], "the earlier call first");
second.settle(again.id, Ok(&[])).expect("settle");
assert_eq!(
wakes(&outcome),
1,
"the settlement by the second handler wakes the caller"
);
let then = second
.next_claim(&mut buf)
.expect("next_claim")
.expect("returned");
assert_eq!(&buf[..then.len], &[2]);
second.settle(then.id, Ok(&[])).expect("settle");
assert_eq!(caller.ack(c), Some(Ok(())));
assert_eq!(caller.ack(later), Some(Ok(())));
}
fn sends_until_busy(caller: &mut impl Caller) -> usize {
for sent in 0..=Loopback::SLOTS {
match caller.command(IFACE, ORD, &[9], None) {
Ok(_) => {}
Err(SendError::Busy) => return sent,
Err(error) => panic!("a send failed other than busy: {error:?}"),
}
}
Loopback::SLOTS + 1
}
fn offer(handler: &mut impl Handler) -> ClaimId {
let mut short = [0u8; 1];
match handler.next_claim(&mut short) {
Err(ReadError::ShortClaim { claim, .. }) => claim,
other => panic!("a buffer shorter than the arguments reports ShortClaim: {other:?}"),
}
}
#[test]
fn a_forgotten_offered_call_is_presented_again_under_the_same_id_and_its_settlement_frees_the_slot()
{
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
handler.serve(IFACE, &[ORD]).expect("serve");
let c = caller.command(IFACE, ORD, &[1, 2, 3], None).expect("send");
let claim = offer(&mut handler);
caller.forget(c);
let mut buf = [0u8; 8];
let retry = handler
.next_claim(&mut buf)
.expect("read")
.expect("a forgotten offered call stays presentable");
assert_eq!(retry.id, claim, "under the id ShortClaim carried");
assert_eq!(&buf[..retry.len], &[1, 2, 3]);
handler.settle(retry.id, Ok(&[])).expect("settle");
assert_eq!(
sends_until_busy(&mut caller),
Loopback::SLOTS,
"the settlement freed the slot"
);
}
#[test]
fn a_dropped_handlers_offered_claim_is_taken_once_by_another_handler_under_the_same_id() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
first.serve(IFACE, &[ORD]).expect("serve");
second.serve(IFACE, &[ORD]).expect("serve");
let c = caller.command(IFACE, ORD, &[1, 2, 3], None).expect("send");
let claim = offer(&mut first);
let (count, waker) = counting();
second.wake_on(Interest::Claim(IFACE), &waker);
let before = wakes(&count);
drop(first);
assert_eq!(
wakes(&count),
before,
"the drop wakes no handler for an offered claim: the call never left the waiting calls"
);
let mut buf = [0u8; 8];
let taken = second
.next_claim(&mut buf)
.expect("read")
.expect("the offered call is still waiting");
assert_eq!(taken.id, claim, "the id stays on the call");
assert_eq!(
second.next_claim(&mut buf).expect("read"),
None,
"the call is presented once, not re-inserted"
);
second.settle(taken.id, Ok(&[])).expect("settle");
assert_eq!(
caller.ack(c),
Some(Ok(())),
"the settlement reaches the caller"
);
}
#[test]
fn a_forgotten_offered_call_is_withdrawn_when_its_handler_drops() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
first.serve(IFACE, &[ORD]).expect("serve");
second.serve(IFACE, &[ORD]).expect("serve");
let c = caller.command(IFACE, ORD, &[1, 2, 3], None).expect("send");
offer(&mut first);
caller.forget(c);
drop(first);
let mut buf = [0u8; 8];
assert_eq!(
second.next_claim(&mut buf).expect("read"),
None,
"a forgotten offered call is withdrawn at the drop"
);
assert_eq!(
sends_until_busy(&mut caller),
Loopback::SLOTS,
"and its slot is reclaimed"
);
}
#[test]
fn an_offered_call_taken_by_another_handler_moves_the_claim_to_it() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
first.serve(IFACE, &[ORD]).expect("serve");
second.serve(IFACE, &[ORD]).expect("serve");
let c = caller.command(IFACE, ORD, &[1, 2, 3], None).expect("send");
let claim = offer(&mut first);
let mut buf = [0u8; 8];
let taken = second
.next_claim(&mut buf)
.expect("read")
.expect("an offered call can be taken by another serving handler");
assert_eq!(taken.id, claim);
second
.settle(taken.id, Ok(&[]))
.expect("the taker settles it");
assert_eq!(
first.settle(claim, Ok(&[])),
Err(SettleError::UnknownClaim),
"the claim moved to the taker"
);
assert_eq!(caller.ack(c), Some(Ok(())));
}
#[test]
fn a_send_wakes_the_handler_that_serves_the_member() {
let rt = runtime();
let mut caller = rt.caller();
let mut serving = rt.handler();
let mut elsewhere = rt.handler();
serving.serve(IFACE, &[ORD]).expect("serve");
elsewhere.serve(IFACE, &[OTHER]).expect("serve");
let (woken, woken_waker) = counting();
let (other, other_waker) = counting();
serving.wake_on(Interest::Claim(IFACE), &woken_waker);
elsewhere.wake_on(Interest::Claim(IFACE), &other_waker);
caller.command(IFACE, ORD, &[1], None).expect("send");
assert_eq!(wakes(&woken), 1, "the serving handler is woken");
assert_eq!(wakes(&other), 0, "a handler serving another member is not");
let mut buf = [0u8; 8];
assert!(serving.next_claim(&mut buf).expect("next_claim").is_some());
}
#[test]
fn a_serve_that_admits_a_waiting_call_wakes_the_handler() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
handler.serve(IFACE, &[OTHER]).expect("serve");
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
caller.command(IFACE, ORD, &[1], None).expect("send");
assert_eq!(wakes(&count), 0, "the handler does not serve the member");
handler.serve(IFACE, &[ORD]).expect("serve");
assert_eq!(wakes(&count), 1, "the call is now waiting for this handler");
}
#[test]
fn a_caller_handle_reads_the_clock() {
let mut rt = runtime();
let caller = rt.caller();
rt.advance(Duration(250));
assert_eq!(caller.now(), Timestamp(250));
assert_eq!(caller.now(), rt.now());
}
#[test]
fn a_registration_whose_key_already_holds_is_woken_at_once() {
let rt = runtime();
let mut caller = rt.caller();
let mut sink = rt.sink();
let mut source = rt.source();
let handler = rt.handler();
source.subscribe(IFACE, &[ORD]).expect("subscribe");
let (slot, slot_waker) = counting();
caller.wake_on(Interest::Slot, &slot_waker);
assert_eq!(wakes(&slot), 1, "a slot is free");
sink.raise(IFACE, ORD, &[1], None).expect("raise");
let (event, event_waker) = counting();
source.wake_on(Interest::Event(IFACE), &event_waker);
assert_eq!(wakes(&event), 1, "an occurrence is waiting");
caller.command(IFACE, ORD, &[1], None).expect("send");
let (claim, claim_waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &claim_waker);
assert_eq!(wakes(&claim), 1, "a call is waiting");
let (unknown, unknown_waker) = counting();
caller.wake_on(
Interest::Outcome(ridl_rt::port::Correlation(99)),
&unknown_waker,
);
assert_eq!(wakes(&unknown), 1, "no call has that correlation");
}
#[test]
fn a_forgotten_call_wakes_its_waiter() {
let Handles { mut caller, .. } = runtime().split();
let c = caller.command(IFACE, ORD, &[1], None).expect("send");
let (count, waker) = counting();
caller.wake_on(Interest::Outcome(c), &waker);
caller.forget(c);
assert_eq!(wakes(&count), 1, "the forget wakes the waiter");
}
#[test]
fn a_key_no_role_of_the_handle_observes_is_woken_at_once() {
let Handles {
reader,
writer,
sink,
caller,
source,
handler,
} = runtime().split();
let (count, waker) = counting();
reader.wake_on(Interest::Slot, &waker);
writer.wake_on(Interest::Event(IFACE), &waker);
sink.wake_on(Interest::Claim(IFACE), &waker);
caller.wake_on(Interest::Event(IFACE), &waker);
caller.wake_on(Interest::Claim(IFACE), &waker);
source.wake_on(Interest::Slot, &waker);
source.wake_on(Interest::Claim(IFACE), &waker);
source.wake_on(Interest::Outcome(ridl_rt::port::Correlation(0)), &waker);
handler.wake_on(Interest::Slot, &waker);
handler.wake_on(Interest::Event(IFACE), &waker);
handler.wake_on(Interest::Outcome(ridl_rt::port::Correlation(0)), &waker);
assert_eq!(wakes(&count), 11);
}
#[test]
fn a_dropped_handler_leaves_no_waiter_behind() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
caller.command(IFACE, ORD, &[1], None).expect("send");
handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
assert_eq!(Arc::strong_count(&count), 3, "the store holds a clone");
drop(handler);
assert_eq!(Arc::strong_count(&count), 2, "the drop released it");
assert_eq!(wakes(&count), 0, "the returned claim does not wake it");
caller.command(IFACE, ORD, &[2], None).expect("send");
assert_eq!(wakes(&count), 0, "nothing wakes a dropped handler's waker");
}
#[test]
fn a_dropped_source_leaves_no_waiter_behind() {
let rt = runtime();
let mut sink = rt.sink();
let mut source = rt.source();
source.subscribe(IFACE, &[ORD]).expect("subscribe");
let (count, waker) = counting();
source.wake_on(Interest::Event(IFACE), &waker);
assert_eq!(Arc::strong_count(&count), 3, "the store holds a clone");
drop(source);
assert_eq!(Arc::strong_count(&count), 2, "the drop released it");
sink.raise(IFACE, ORD, &[1], None).expect("raise");
assert_eq!(wakes(&count), 0, "nothing wakes a dropped source's waker");
}
#[test]
fn a_providers_busy_settlement_reaches_the_caller() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
let busy = CallError::Transport(Transport::Busy);
let command = caller.command(IFACE, ORD, &[1], None).expect("send");
let query = caller.query(IFACE, ORD, &[2], None).expect("send");
while let Some(claim) = handler.next_claim(&mut buf).expect("next_claim") {
handler.settle(claim.id, Err(busy)).expect("settle");
}
assert_eq!(caller.ack(command), Some(Err(busy)));
assert_eq!(caller.reply(query, &mut buf), Ok(Some(Err(busy))));
}
#[test]
fn a_returned_claim_keeps_its_place_by_send_order() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
let mut buf = [0u8; 8];
caller.command(IFACE, ORD, &[1], None).expect("send");
first
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
caller.command(IFACE, ORD, &[2], None).expect("send");
drop(first);
let claim = second
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
&buf[..claim.len],
&[1],
"the returned call was sent first, so it is presented first"
);
}
#[test]
fn a_returned_claim_keeps_its_send_order_across_a_reused_slot() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
let mut buf = [0u8; 8];
let old = caller.command(IFACE, ORD, &[0], None).expect("send");
let claim = first
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
first.settle(claim.id, Ok(&[])).expect("settle");
caller.forget(old);
let reused = caller
.command(IFACE, ORD, &[1], None)
.expect("send into slot 0");
let fresh = caller
.command(IFACE, ORD, &[2], None)
.expect("send into slot 1");
assert!(
reused.0 > fresh.0,
"the reused slot's correlation is the larger one"
);
first
.next_claim(&mut buf)
.expect("next_claim")
.expect("the call sent first");
drop(first);
let claim = second
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
&buf[..claim.len],
&[1],
"the returned call was sent first, so it is presented first"
);
}
fn fill(caller: &mut impl Caller) -> Vec<ridl_rt::port::Correlation> {
(0..Loopback::SLOTS)
.map(|n| {
caller
.command(IFACE, ORD, &[u8::try_from(n).expect("small")], None)
.expect("a slot is free")
})
.collect()
}
#[test]
fn the_seventeenth_in_flight_call_is_busy() {
assert_eq!(Loopback::SLOTS, 16);
let rt = runtime();
let mut caller = rt.caller();
let mut other = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
let calls = fill(&mut caller);
assert_eq!(
caller.command(IFACE, ORD, &[99], None),
Err(SendError::Busy)
);
assert_eq!(caller.query(IFACE, ORD, &[99], None), Err(SendError::Busy));
assert_eq!(
other.command(IFACE, ORD, &[99], None),
Err(SendError::Busy),
"the table is the runtime's, shared by every caller"
);
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
handler.settle(claim.id, Ok(&[])).expect("settle");
assert_eq!(caller.ack(calls[0]), Some(Ok(())));
assert_eq!(
caller.command(IFACE, ORD, &[99], None),
Err(SendError::Busy),
"a settled call keeps its slot until it is forgotten"
);
caller.forget(calls[0]);
other
.command(IFACE, ORD, &[99], None)
.expect("the forget freed a slot");
assert_eq!(
caller.command(IFACE, ORD, &[100], None),
Err(SendError::Busy)
);
}
#[test]
fn a_refused_send_draws_no_sequence_number() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
let calls = fill(&mut caller);
assert_eq!(
caller.command(IFACE, ORD, &[99], None),
Err(SendError::Busy)
);
let mut last = 0;
for _ in &calls {
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
last = claim.envelope.seq;
handler.settle(claim.id, Ok(&[])).expect("settle");
}
caller.forget(calls[0]);
caller
.command(IFACE, ORD, &[99], None)
.expect("a slot is free");
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
claim.envelope.seq,
last + 1,
"nothing was sent, so no number was used"
);
}
#[test]
fn a_claimed_then_forgotten_call_holds_its_slot_until_its_settlement() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
let calls = fill(&mut caller);
let (count, waker) = counting();
caller.wake_on(Interest::Slot, &waker);
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(&buf[..claim.len], &[0], "the claim is calls[0]");
caller.forget(calls[0]);
assert_eq!(wakes(&count), 0, "the claimed call still holds its slot");
assert_eq!(
caller.command(IFACE, ORD, &[99], None),
Err(SendError::Busy)
);
handler
.settle(claim.id, Ok(&[]))
.expect("the provider's settlement is still accepted");
assert_eq!(wakes(&count), 1, "its settlement reclaims the slot");
caller
.command(IFACE, ORD, &[99], None)
.expect("the slot is free");
}
#[test]
fn forgotten_calls_to_an_unserved_member_leave_room_for_another_send() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
handler.serve(IFACE, &[ORD]).expect("serve");
let calls: Vec<_> = (0..Loopback::SLOTS)
.map(|n| {
caller
.command(IFACE, OTHER, &[u8::try_from(n).expect("small")], None)
.expect("a slot is free")
})
.collect();
assert_eq!(
caller.command(IFACE, OTHER, &[99], None),
Err(SendError::Busy)
);
let (count, waker) = counting();
caller.wake_on(Interest::Slot, &waker);
for c in &calls {
caller.forget(*c);
}
assert_eq!(wakes(&count), 1, "the withdrawal reclaimed a slot");
for n in 0..Loopback::SLOTS {
caller
.command(IFACE, OTHER, &[u8::try_from(n).expect("small")], None)
.expect("every withdrawn call's slot is free");
}
assert_eq!(
caller.command(IFACE, OTHER, &[99], None),
Err(SendError::Busy)
);
}
#[test]
fn a_withdrawn_call_is_never_presented_to_a_handler_that_serves_the_member_later() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
handler.serve(IFACE, &[OTHER]).expect("serve");
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
let withdrawn = caller.command(IFACE, ORD, &[1], None).expect("send");
let kept = caller.command(IFACE, ORD, &[2], None).expect("send");
caller.forget(withdrawn);
handler.serve(IFACE, &[ORD]).expect("serve");
assert_eq!(wakes(&count), 1, "the call not forgotten is waiting");
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("the call not forgotten");
assert_eq!(&buf[..claim.len], &[2], "the withdrawn call is skipped");
handler.settle(claim.id, Ok(&[])).expect("settle");
assert_eq!(caller.ack(kept), Some(Ok(())));
assert_eq!(
handler.next_claim(&mut buf).expect("next_claim"),
None,
"the withdrawn call is never presented"
);
}
#[test]
fn a_withdrawal_from_the_middle_of_the_queue_leaves_the_calls_around_it_in_order() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
caller.command(IFACE, ORD, &[1], None).expect("send");
let middle = caller.command(IFACE, ORD, &[2], None).expect("send");
caller.command(IFACE, ORD, &[3], None).expect("send");
caller.forget(middle);
handler.serve(IFACE, &[ORD]).expect("serve");
for expected in [1u8, 3] {
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
&buf[..claim.len],
&[expected],
"send order, less the middle"
);
}
assert_eq!(handler.next_claim(&mut buf).expect("next_claim"), None);
}
#[test]
fn forgetting_a_claimed_call_wakes_its_outcome_waiter() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let c = caller.command(IFACE, ORD, &[1], None).expect("send");
handler
.next_claim(&mut [0u8; 8])
.expect("next_claim")
.expect("waiting");
let (count, waker) = counting();
caller.wake_on(Interest::Outcome(c), &waker);
assert_eq!(wakes(&count), 0, "the claimed call is in flight");
caller.forget(c);
assert_eq!(
wakes(&count),
1,
"no outcome will be readable for a forgotten call"
);
}
#[test]
fn a_stale_correlation_does_not_withdraw_the_call_now_in_its_slot() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
let old = caller.command(IFACE, ORD, &[1], None).expect("send");
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
handler.settle(claim.id, Ok(&[])).expect("settle");
caller.forget(old);
let new = caller
.command(IFACE, ORD, &[2], None)
.expect("send into the same slot");
assert_ne!(new, old, "the slot is reused under a new generation");
caller.forget(old);
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("the stale correlation withdrew nothing");
assert_eq!(&buf[..claim.len], &[2]);
handler.settle(claim.id, Ok(&[])).expect("settle");
assert_eq!(caller.ack(new), Some(Ok(())));
}
#[test]
fn a_call_forgotten_while_claimed_is_withdrawn_when_its_handler_is_dropped() {
let rt = runtime();
let mut caller = rt.caller();
let mut buf = [0u8; 8];
let mut handlers = Vec::new();
for n in 0..Loopback::SLOTS {
let c = caller
.command(IFACE, ORD, &[u8::try_from(n).expect("small")], None)
.expect("a slot is free");
let mut handler = rt.handler();
handler.serve(IFACE, &[ORD]).expect("serve");
handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
caller.forget(c);
handlers.push(handler);
}
assert_eq!(
caller.command(IFACE, ORD, &[99], None),
Err(SendError::Busy),
"a claimed call keeps its slot after its forget"
);
let (count, waker) = counting();
caller.wake_on(Interest::Slot, &waker);
assert_eq!(wakes(&count), 0, "no slot is free: the waker is stored");
drop(handlers);
assert_eq!(wakes(&count), 1, "the drops reclaimed the slots");
for n in 0..Loopback::SLOTS {
caller
.command(IFACE, ORD, &[100 + u8::try_from(n).expect("small")], None)
.expect("every withdrawn call's slot is free");
}
assert_eq!(
caller.command(IFACE, ORD, &[99], None),
Err(SendError::Busy)
);
let mut later = rt.handler();
later.serve(IFACE, &[ORD]).expect("serve");
for n in 0..Loopback::SLOTS {
let claim = later
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
&buf[..claim.len],
&[100 + u8::try_from(n).expect("small")],
"no withdrawn call is presented again"
);
}
assert_eq!(later.next_claim(&mut buf).expect("next_claim"), None);
}
#[test]
fn a_forget_withdrawal_wakes_no_claim_waiter_of_a_handler_serving_another_member() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
handler.serve(IFACE, &[OTHER]).expect("serve");
let c = caller.command(IFACE, ORD, &[1], None).expect("send");
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
assert_eq!(wakes(&count), 0, "no call it serves waits: stored");
caller.forget(c);
assert_eq!(
wakes(&count),
0,
"a withdrawal adds no call to claim, so it wakes no claim waiter"
);
caller
.command(IFACE, OTHER, &[2], None)
.expect("send a call it serves");
assert_eq!(wakes(&count), 1, "the waker was still stored");
}
#[test]
fn a_dropped_handler_withdraws_only_its_forgotten_claim_and_returns_the_others() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut buf = [0u8; 8];
first.serve(IFACE, &[ORD]).expect("serve");
let sent: Vec<_> = (1u8..=3)
.map(|n| caller.command(IFACE, ORD, &[n], None).expect("send"))
.collect();
for _ in &sent {
first
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
}
for n in 3..Loopback::SLOTS {
caller
.command(IFACE, OTHER, &[u8::try_from(n).expect("small")], None)
.expect("a slot is free");
}
caller.forget(sent[1]);
assert_eq!(
caller.command(IFACE, OTHER, &[99], None),
Err(SendError::Busy),
"the claimed call keeps its slot after its forget"
);
drop(first);
caller
.command(IFACE, OTHER, &[99], None)
.expect("the withdrawn claim's slot came back");
assert_eq!(
caller.command(IFACE, OTHER, &[100], None),
Err(SendError::Busy),
"exactly one slot came back"
);
let mut later = rt.handler();
later.serve(IFACE, &[ORD]).expect("serve");
for expected in [1u8, 3] {
let claim = later
.next_claim(&mut buf)
.expect("next_claim")
.expect("a returned claim");
assert_eq!(&buf[..claim.len], &[expected], "send order, less 2");
}
assert_eq!(later.next_claim(&mut buf).expect("next_claim"), None);
}
#[test]
fn a_dropped_callers_claimed_call_is_withdrawn_when_its_handler_is_dropped() {
let rt = runtime();
let mut caller = rt.caller();
let mut other = rt.caller();
let mut first = rt.handler();
let mut buf = [0u8; 8];
first.serve(IFACE, &[ORD]).expect("serve");
caller.command(IFACE, ORD, &[1], None).expect("send");
first
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
drop(caller);
drop(first);
for n in 0..Loopback::SLOTS {
other
.command(IFACE, ORD, &[100 + u8::try_from(n).expect("small")], None)
.expect("the withdrawn call's slot is free");
}
let mut later = rt.handler();
later.serve(IFACE, &[ORD]).expect("serve");
let claim = later
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
&buf[..claim.len],
&[100],
"the dropped caller's call is not presented again"
);
}
#[test]
fn a_withdrawal_at_a_handlers_drop_wakes_no_claim_waiter() {
let rt = runtime();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
first.serve(IFACE, &[ORD]).expect("serve");
second.serve(IFACE, &[ORD]).expect("serve");
let c = caller.command(IFACE, ORD, &[1], None).expect("send");
first
.next_claim(&mut [0u8; 8])
.expect("next_claim")
.expect("waiting");
let (count, waker) = counting();
second.wake_on(Interest::Claim(IFACE), &waker);
assert_eq!(wakes(&count), 0, "nothing is waiting: the waker is stored");
caller.forget(c);
drop(first);
assert_eq!(
wakes(&count),
0,
"a withdrawal adds no call to claim, so it wakes no claim waiter"
);
assert_eq!(second.next_claim(&mut [0u8; 8]).expect("next_claim"), None);
}
#[test]
fn a_slot_registration_while_a_slot_is_free_is_woken_at_once() {
let rt = runtime();
let mut caller = rt.caller();
let calls = fill(&mut caller);
let mut handler = rt.handler();
let mut buf = [0u8; 8];
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
handler.settle(claim.id, Ok(&[])).expect("settle");
caller.forget(calls[0]);
let (count, waker) = counting();
caller.wake_on(Interest::Slot, &waker);
assert_eq!(wakes(&count), 1, "one slot is free");
caller.command(IFACE, ORD, &[99], None).expect("send");
caller.wake_on(Interest::Slot, &waker);
assert_eq!(wakes(&count), 1, "the table is full again: stored");
}
#[test]
fn a_dropped_caller_leaves_no_slot_waiter_behind() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
let calls = fill(&mut caller);
let other = rt.caller();
let (count, waker) = counting();
other.wake_on(Interest::Slot, &waker);
assert_eq!(Arc::strong_count(&count), 3, "the store holds a clone");
drop(other);
assert_eq!(Arc::strong_count(&count), 2, "the drop released it");
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
handler.settle(claim.id, Ok(&[])).expect("settle");
caller.forget(calls[0]);
assert_eq!(wakes(&count), 0);
}
#[test]
fn a_dropped_caller_forgets_its_calls_and_their_slots_come_back() {
let rt = runtime();
let mut caller = rt.caller();
let mut other = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
fill(&mut caller);
for _ in 0..8 {
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
handler.settle(claim.id, Ok(&[])).expect("settle");
}
let held: Vec<_> = (0..8)
.map(|_| {
handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting")
})
.collect();
assert_eq!(other.command(IFACE, ORD, &[99], None), Err(SendError::Busy));
let (count, waker) = counting();
other.wake_on(Interest::Slot, &waker);
drop(caller);
assert_eq!(
wakes(&count),
1,
"the drop reclaimed the settled calls' slots"
);
for n in 0..8 {
other
.command(IFACE, ORD, &[n], None)
.expect("a slot of a settled call of the dropped caller");
}
assert_eq!(
other.command(IFACE, ORD, &[99], None),
Err(SendError::Busy),
"the dropped caller's claimed calls keep their slots until settled"
);
for claim in held {
handler.settle(claim.id, Ok(&[])).expect("settle");
}
for n in 0..8 {
other
.command(IFACE, ORD, &[n], None)
.expect("a slot of a claimed call of the dropped caller");
}
assert_eq!(other.command(IFACE, ORD, &[99], None), Err(SendError::Busy));
}
#[test]
fn a_dropped_callers_unclaimed_calls_are_withdrawn() {
let rt = runtime();
let mut caller = rt.caller();
let mut other = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
fill(&mut caller);
let (count, waker) = counting();
other.wake_on(Interest::Slot, &waker);
drop(caller);
assert_eq!(wakes(&count), 1, "the drop reclaimed the unclaimed calls");
for n in 0..Loopback::SLOTS {
other
.command(IFACE, ORD, &[100 + u8::try_from(n).expect("small")], None)
.expect("every unclaimed call of the dropped caller was withdrawn");
}
assert_eq!(other.command(IFACE, ORD, &[99], None), Err(SendError::Busy));
for n in 0..Loopback::SLOTS {
let claim = handler
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
assert_eq!(
&buf[..claim.len],
&[100 + u8::try_from(n).expect("small")],
"only the other caller's calls are presented"
);
}
assert_eq!(handler.next_claim(&mut buf).expect("next_claim"), None);
}
#[test]
fn a_dropped_caller_forgets_only_its_own_calls() {
let rt = runtime();
let mut first = rt.caller();
let mut second = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
first.command(IFACE, ORD, &[1], None).expect("send");
let theirs = second.command(IFACE, ORD, &[2], None).expect("send");
while let Some(claim) = handler.next_claim(&mut buf).expect("next_claim") {
handler.settle(claim.id, Ok(&[])).expect("settle");
}
drop(first);
assert_eq!(
second.ack(theirs),
Some(Ok(())),
"the other caller's settled call is still readable"
);
}
#[test]
fn a_dropped_callers_call_in_flight_wakes_its_outcome_waiter() {
let rt = runtime();
let mut caller = rt.caller();
let c = caller.command(IFACE, ORD, &[1], None).expect("send");
let (count, waker) = counting();
caller.wake_on(Interest::Outcome(c), &waker);
assert_eq!(wakes(&count), 0, "the call is in flight");
drop(caller);
assert_eq!(
wakes(&count),
1,
"the drop forgot the call, so no outcome will be readable for it"
);
}
#[test]
fn a_dropped_callers_own_slot_waiter_is_not_woken_by_its_drop() {
let rt = runtime();
let mut caller = rt.caller();
let other = rt.caller();
let mut handler = rt.handler();
let mut buf = [0u8; 8];
fill(&mut caller);
while let Some(claim) = handler.next_claim(&mut buf).expect("next_claim") {
handler.settle(claim.id, Ok(&[])).expect("settle");
}
let (own, own_waker) = counting();
let (theirs, their_waker) = counting();
caller.wake_on(Interest::Slot, &own_waker);
other.wake_on(Interest::Slot, &their_waker);
drop(caller);
assert_eq!(
(wakes(&own), wakes(&theirs)),
(0, 1),
"the dropped caller's waiters leave before its calls are reclaimed"
);
assert_eq!(Arc::strong_count(&own), 2, "and its waker is released");
}
#[test]
fn a_serve_with_no_call_waiting_keeps_the_handlers_claim_waker() {
let rt = runtime();
let mut caller = rt.caller();
let mut handler = rt.handler();
let (count, waker) = counting();
handler.wake_on(Interest::Claim(IFACE), &waker);
handler.serve(IFACE, &[ORD]).expect("serve");
assert_eq!(wakes(&count), 0, "nothing is waiting");
caller.command(IFACE, ORD, &[1], None).expect("send");
assert_eq!(
wakes(&count),
1,
"the waker is still stored, so the send wakes it"
);
}
struct OwnsAHandle<H>(std::sync::Mutex<Option<H>>);
impl<H: Send + 'static> Wake for OwnsAHandle<H> {
fn wake(self: Arc<Self>) {
let owned = self
.0
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
assert!(
owned.is_some(),
"the handle is owned until the waker is dropped"
);
}
}
fn owning<H: Send + 'static>(handle: H) -> Waker {
Waker::from(Arc::new(OwnsAHandle(std::sync::Mutex::new(Some(handle)))))
}
fn assert_drop_finishes<H: Send + 'static>(handle: H, what: &str) {
let (done, finished) = mpsc::channel();
std::thread::spawn(move || {
drop(handle);
let _ = done.send(());
});
assert_eq!(
finished.recv_timeout(std::time::Duration::from_secs(5)),
Ok(()),
"{what}: the drop finished, so no waker was dropped under the store's lock"
);
}
#[test]
fn a_caller_whose_stored_waker_owns_another_handle_drops_without_deadlock() {
let rt = runtime();
let mut caller = rt.caller();
fill(&mut caller);
let waker = owning(rt.caller());
caller.wake_on(Interest::Slot, &waker);
drop(waker);
std::mem::forget(rt);
assert_drop_finishes(caller, "a caller");
}
#[test]
fn a_source_whose_stored_waker_owns_another_handle_drops_without_deadlock() {
let rt = runtime();
let source = rt.source();
let waker = owning(rt.source());
source.wake_on(Interest::Event(IFACE), &waker);
drop(waker);
std::mem::forget(rt);
assert_drop_finishes(source, "a source");
}
#[test]
fn a_handler_whose_stored_waker_owns_another_handle_drops_without_deadlock() {
let rt = runtime();
let handler = rt.handler();
let waker = owning(rt.handler());
handler.wake_on(Interest::Claim(IFACE), &waker);
drop(waker);
std::mem::forget(rt);
assert_drop_finishes(handler, "a handler");
}
#[test]
fn the_aggregate_routes_each_key_to_the_handle_that_observes_it() {
let mut rt = runtime();
let mut buf = [0u8; 8];
let c = rt.command(IFACE, ORD, &[1], None).expect("send");
let (outcome, outcome_waker) = counting();
rt.wake_on(Interest::Outcome(c), &outcome_waker);
assert_eq!(wakes(&outcome), 0, "stored, not woken at once");
let claim = rt
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
rt.settle(claim.id, Ok(&[])).expect("settle");
assert_eq!(wakes(&outcome), 1, "woken by the settlement");
rt.subscribe(IFACE, &[ORD]).expect("subscribe");
let (event, event_waker) = counting();
rt.wake_on(Interest::Event(IFACE), &event_waker);
assert_eq!(wakes(&event), 0, "stored, not woken at once");
rt.raise(IFACE, ORD, &[1], None).expect("raise");
assert_eq!(wakes(&event), 1, "woken by the raise");
let (claim, claim_waker) = counting();
rt.wake_on(Interest::Claim(IFACE), &claim_waker);
assert_eq!(wakes(&claim), 0, "stored, not woken at once");
rt.caller().command(IFACE, ORD, &[1], None).expect("send");
assert_eq!(wakes(&claim), 1, "woken by the send");
let (slot, slot_waker) = counting();
rt.wake_on(Interest::Slot, &slot_waker);
assert_eq!(
wakes(&slot),
1,
"the caller wakes it at once: a slot is free"
);
}
struct ReadsTheStore {
reader: Arc<ridl_loopback::ReaderHandle>,
done: std::sync::Mutex<Option<mpsc::Sender<bool>>>,
}
impl Wake for ReadsTheStore {
fn wake(self: Arc<Self>) {
let (tx, rx) = mpsc::channel();
let reader = Arc::clone(&self.reader);
std::thread::spawn(move || {
reader.now();
let _ = tx.send(());
});
let finished = rx.recv_timeout(std::time::Duration::from_secs(1)).is_ok();
if let Some(done) = self.done.lock().expect("not poisoned").take() {
let _ = done.send(finished);
}
}
}
fn lock_probe(rt: &Loopback) -> (Waker, mpsc::Receiver<bool>) {
let (tx, rx) = mpsc::channel();
let waker = Waker::from(Arc::new(ReadsTheStore {
reader: Arc::new(rt.reader()),
done: std::sync::Mutex::new(Some(tx)),
}));
(waker, rx)
}
fn assert_released(rx: &mpsc::Receiver<bool>, path: &str) {
assert_eq!(
rx.recv_timeout(std::time::Duration::from_secs(5)),
Ok(true),
"{path}: the waker ran with the store's lock released"
);
}
#[test]
fn every_wake_is_run_with_the_lock_released() {
let rt = runtime();
let mut sink = rt.sink();
let mut source = rt.source();
let mut caller = rt.caller();
let mut first = rt.handler();
let mut second = rt.handler();
source.subscribe(IFACE, &[ORD]).expect("subscribe");
first.serve(IFACE, &[ORD]).expect("serve");
second.serve(IFACE, &[OTHER]).expect("serve");
let mut buf = [0u8; 8];
let (waker, rx) = lock_probe(&rt);
source.wake_on(Interest::Event(IFACE), &waker);
sink.raise(IFACE, ORD, &[1], None).expect("raise");
assert_released(&rx, "a raise");
let (waker, rx) = lock_probe(&rt);
source.wake_on(Interest::Event(IFACE), &waker);
assert_released(&rx, "a registration whose key already holds");
let (waker, rx) = lock_probe(&rt);
first.wake_on(Interest::Claim(IFACE), &waker);
caller.command(IFACE, ORD, &[1], None).expect("send");
assert_released(&rx, "a command");
first
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
let (waker, rx) = lock_probe(&rt);
first.wake_on(Interest::Claim(IFACE), &waker);
let c = caller.query(IFACE, ORD, &[1], None).expect("send");
assert_released(&rx, "a query");
let (waker, rx) = lock_probe(&rt);
caller.wake_on(Interest::Outcome(c), &waker);
caller.forget(c);
assert_released(&rx, "a forget");
caller.query(IFACE, ORD, &[1], None).expect("send");
let (waker, rx) = lock_probe(&rt);
second.wake_on(Interest::Claim(IFACE), &waker);
second.serve(IFACE, &[ORD]).expect("serve");
assert_released(&rx, "a serve");
second
.next_claim(&mut buf)
.expect("next_claim")
.expect("waiting");
let (waker, rx) = lock_probe(&rt);
first.wake_on(Interest::Claim(IFACE), &waker);
assert!(first.next_claim(&mut buf).expect("next_claim").is_none());
drop(second);
assert_released(&rx, "a handler's drop");
let claim = first
.next_claim(&mut buf)
.expect("next_claim")
.expect("the returned query");
first.settle(claim.id, Ok(&[])).expect("settle");
let mut sent = Vec::new();
while let Ok(c) = caller.command(IFACE, ORD, &[2], None) {
sent.push(c);
}
let claim = first
.next_claim(&mut buf)
.expect("next_claim")
.expect("the first command");
first.settle(claim.id, Ok(&[])).expect("settle");
let (waker, rx) = lock_probe(&rt);
caller.wake_on(Interest::Slot, &waker);
caller.forget(sent[0]);
assert_released(&rx, "a forget that reclaims a slot");
caller
.command(IFACE, ORD, &[3], None)
.expect("the reclaimed slot");
let claim = first
.next_claim(&mut buf)
.expect("next_claim")
.expect("the second command");
caller.forget(sent[1]);
let (waker, rx) = lock_probe(&rt);
caller.wake_on(Interest::Slot, &waker);
first.settle(claim.id, Ok(&[])).expect("settle");
assert_released(&rx, "a settlement that reclaims a slot");
caller
.command(IFACE, ORD, &[4], None)
.expect("the reclaimed slot");
let (waker, rx) = lock_probe(&rt);
caller.wake_on(Interest::Slot, &waker);
caller.forget(sent[2]);
assert_released(&rx, "a forget that withdraws an unclaimed call");
}