use super::*;
use crate::session::{MINIMUM_COMPLETION_RING_CAPACITY, MINIMUM_SUBMISSION_CAPACITY};
#[test]
fn one_enumeration_runs_to_completion() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::OfferEntry(0, "a"),
Op::OfferEntry(0, "b"),
Op::Report(Quantum::Idle),
Op::Schedule(0),
Op::RunEngine(Quantum::Completed),
Op::Service,
Op::Detach(0),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a", "b"]);
assert_eq!(model.terminal(0), Some("completed"));
assert_eq!(model.registered(), 0);
}
#[test]
fn two_enumerations_interleave_without_losing_their_own_order() {
let mut model = Model::new(16, 16);
model.run(&[
Op::Begin,
Op::Begin,
Op::Service,
Op::OfferEntry(0, "a0"),
Op::OfferEntry(1, "b0"),
Op::OfferEntry(0, "a1"),
Op::OfferEntry(1, "b1"),
Op::OfferEntry(0, "a2"),
Op::RunEngine(Quantum::Completed),
Op::Schedule(1),
Op::RunEngine(Quantum::Completed),
Op::Service,
Op::Detach(0),
Op::Detach(1),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a0", "a1", "a2"]);
assert_eq!(model.entries(1), ["b0", "b1"]);
assert_eq!(model.terminal(0), Some("completed"));
assert_eq!(model.terminal(1), Some("completed"));
}
#[test]
fn cancelling_a_quiescent_enumeration_terminates_it() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Cancel(0),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.terminal(0), Some("cancelled"));
assert_eq!(model.registered(), 0);
}
#[test]
fn dropping_the_handle_terminates_the_enumeration() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::DropHandle(0),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.terminal(0), Some("cancelled"));
}
#[test]
fn cancelling_during_a_quantum_defers_the_terminal_behind_its_entries() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::OfferEntry(0, "a"),
Op::Cancel(0),
Op::Service,
]);
assert_eq!(model.registered(), 1, "the quantum still owns it");
assert_eq!(model.terminal(0), None);
model.run(&[
Op::OfferEntry(0, "b"),
Op::Report(Quantum::Idle),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a", "b"]);
assert_eq!(model.terminal(0), Some("cancelled"));
assert_eq!(model.registered(), 0);
}
#[test]
fn a_quantum_scheduled_after_cancellation_does_nothing() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Cancel(0),
Op::Service,
Op::Claim,
Op::Report(Quantum::Idle),
Op::DrainReceiver,
]);
assert_eq!(model.terminal(0), Some("cancelled"));
}
#[test]
fn a_cancel_for_an_unknown_enumeration_is_ignored() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Cancel(0),
Op::Service,
Op::Begin,
Op::Service,
Op::Detach(1),
]);
assert_eq!(model.registered(), 1);
assert_eq!(model.terminal(1), None);
}
#[test]
fn backpressure_refuses_entries_and_resumes_after_a_take() {
let mut model = Model::new(8, 4);
model.run(&[
Op::Begin,
Op::Service,
Op::OfferEntry(0, "a"),
Op::OfferEntry(0, "b"),
Op::OfferEntry(0, "c"),
Op::OfferEntry(0, "d"),
]);
assert_eq!(model.refused(), 1, "the fourth entry had nowhere to go");
model.run(&[Op::Recv, Op::OfferEntry(0, "d")]);
assert_eq!(model.refused(), 1, "room appeared, so the retry succeeded");
model.run(&[
Op::RunEngine(Quantum::Completed),
Op::Detach(0),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a", "b", "c", "d"]);
assert_eq!(model.terminal(0), Some("completed"));
}
#[test]
fn backpressure_is_shared_between_enumerations_in_one_session() {
let mut model = Model::new(16, 5);
model.run(&[
Op::Begin,
Op::Begin,
Op::Service,
Op::OfferEntry(0, "a0"),
Op::OfferEntry(0, "a1"),
Op::OfferEntry(0, "a2"),
Op::OfferEntry(1, "b0"),
]);
assert_eq!(model.refused(), 1, "the second enumeration is stalled too");
model.run(&[Op::Recv, Op::OfferEntry(1, "b0"), Op::DrainReceiver]);
assert_eq!(model.entries(0), ["a0", "a1", "a2"]);
assert_eq!(model.entries(1), ["b0"]);
model.run(&[Op::Detach(0), Op::Detach(1)]);
}
#[test]
fn a_parked_enumeration_is_resumed_by_a_take() {
let mut model = Model::new(8, 3);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::OfferEntry(0, "a"),
Op::OfferEntry(0, "b"),
Op::OfferEntry(0, "c"),
Op::Report(Quantum::Parked),
]);
assert_eq!(model.refused(), 1);
model.run(&[
Op::Recv,
Op::OfferEntry(0, "c"),
Op::RunEngine(Quantum::Completed),
Op::Service,
Op::Detach(0),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a", "b", "c"]);
assert_eq!(model.terminal(0), Some("completed"));
}
#[test]
fn a_worker_that_parks_after_the_receiver_already_drained_is_not_stranded() {
let mut model = Model::new(8, 3);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::OfferEntry(0, "a"),
Op::OfferEntry(0, "b"),
Op::OfferEntry(0, "c"),
]);
assert_eq!(model.refused(), 1);
assert_eq!(
model.ready(),
0,
"the enumeration is still held, not queued"
);
model.run(&[Op::DrainReceiver]);
assert_eq!(model.entries(0), ["a", "b"]);
model.run(&[Op::Report(Quantum::Parked)]);
assert_eq!(
model.ready(),
1,
"room was already available when the worker reported Parked; it must \
resume itself rather than trust a wakeup that already happened"
);
model.run(&[
Op::RunEngine(Quantum::Completed),
Op::Service,
Op::Detach(0),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a", "b"]);
assert_eq!(model.terminal(0), Some("completed"));
}
#[test]
fn a_terminal_lands_in_a_ring_with_no_ordinary_room() {
let mut model = Model::new(8, 3);
model.run(&[
Op::Begin,
Op::Service,
Op::OfferEntry(0, "a"),
Op::OfferEntry(0, "b"),
Op::OfferEntry(0, "c"),
]);
assert_eq!(model.refused(), 1);
model.run(&[
Op::RunEngine(Quantum::Completed),
Op::Detach(0),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a", "b"]);
assert_eq!(model.terminal(0), Some("completed"));
}
#[test]
fn abandonment_releases_everything_without_a_terminal() {
let mut model = Model::new(16, 16);
model.run(&[
Op::Begin,
Op::Begin,
Op::Service,
Op::OfferEntry(0, "a"),
Op::DropReceiver,
Op::Service,
Op::BeginRefused(BeginFailure::Abandoned),
]);
assert_eq!(model.registered(), 0);
assert_eq!(model.terminal(0), None);
assert_eq!(model.terminal(1), None);
model.run(&[Op::Detach(0), Op::Detach(1)]);
}
#[test]
fn abandonment_during_a_quantum_is_safe() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::DropReceiver,
Op::Service,
Op::OfferEntry(0, "a"),
Op::Report(Quantum::Idle),
Op::RunEngine(Quantum::Completed),
Op::Detach(0),
]);
assert_eq!(model.registered(), 0);
}
#[test]
fn cancelling_a_completed_enumeration_adds_no_terminal() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::RunEngine(Quantum::Completed),
Op::Cancel(0),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.terminal(0), Some("completed"));
assert_eq!(model.registered(), 0);
}
#[test]
fn a_detached_enumeration_still_reports_its_outcome() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Detach(0),
Op::OfferEntry(0, "a"),
Op::RunEngine(Quantum::Completed),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a"]);
assert_eq!(model.terminal(0), Some("completed"));
}
#[test]
fn dropping_the_session_does_not_strand_an_owed_terminal() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::OfferEntry(0, "a"),
Op::RunEngine(Quantum::Completed),
Op::Detach(0),
Op::DropSession,
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a"]);
assert_eq!(model.terminal(0), Some("completed"));
}
#[test]
fn the_minimum_bounds_carry_one_enumeration() {
let mut model = Model::new(
MINIMUM_SUBMISSION_CAPACITY,
MINIMUM_COMPLETION_RING_CAPACITY,
);
model.run(&[
Op::Begin,
Op::BeginRefused(BeginFailure::SubmissionRingFull),
Op::Service,
Op::OfferEntry(0, "only"),
Op::OfferEntry(0, "refused"),
]);
assert_eq!(model.refused(), 1);
model.run(&[
Op::RunEngine(Quantum::Completed),
Op::Detach(0),
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["only"]);
assert_eq!(model.terminal(0), Some("completed"));
}
#[test]
fn a_completion_ring_of_one_is_rejected() {
let error = Session::new(8, 1).expect_err("a ring of one cannot hold both");
assert_eq!(
error.failure(),
crate::error::SessionFailure::CompletionCapacityTooSmall
);
}
#[test]
fn admission_stops_at_the_completion_ring_s_reservation_boundary() {
let mut model = Model::new(32, 3);
model.run(&[
Op::Begin,
Op::Begin,
Op::BeginRefused(BeginFailure::CompletionRingFull),
Op::Service,
]);
assert_eq!(model.registered(), 2);
model.run(&[Op::Detach(0), Op::Detach(1)]);
}
#[test]
fn repeated_cycles_leak_nothing() {
let mut model = Model::new(MINIMUM_SUBMISSION_CAPACITY, 4);
for _ in 0..5 {
model.run(&[Op::Begin, Op::Service]);
let slot = model.ids.len() - 1;
model.run(&[
Op::OfferEntry(slot, "x"),
Op::Cancel(slot),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.entries(slot), ["x"]);
assert_eq!(model.terminal(slot), Some("cancelled"));
assert_eq!(model.registered(), 0);
}
}
#[test]
fn redundant_servicing_changes_nothing() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Service,
Op::Service,
Op::Begin,
Op::Service,
Op::Service,
Op::Service,
Op::Detach(0),
]);
assert_eq!(model.registered(), 1);
}
#[test]
fn draining_an_empty_ring_changes_nothing() {
let mut model = Model::new(8, 8);
model.run(&[
Op::DrainReceiver,
Op::Recv,
Op::Begin,
Op::Service,
Op::DrainReceiver,
Op::Detach(0),
]);
assert_eq!(model.registered(), 1);
}
#[test]
fn a_worker_reports_and_the_servicer_retires() {
let mut model = Model::new(8, 8);
model.run(&[Op::Begin, Op::Service, Op::RunEngine(Quantum::Completed)]);
assert_eq!(
model.registered(),
1,
"the entry survives until the report is serviced"
);
model.run(&[Op::DrainReceiver]);
assert_eq!(
model.terminal(0),
Some("completed"),
"the terminal is the worker's to deliver"
);
model.run(&[Op::Service]);
assert_eq!(model.registered(), 0);
model.run(&[Op::Detach(0)]);
}
#[test]
fn a_retire_serviced_after_abandonment_is_a_no_op() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::RunEngine(Quantum::Completed),
Op::DropReceiver,
Op::Service,
]);
assert_eq!(model.registered(), 0);
model.run(&[Op::Service, Op::Detach(0)]);
assert_eq!(model.registered(), 0);
}
#[test]
fn a_report_after_abandonment_finds_its_enumeration_gone() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::DropReceiver,
Op::Service,
Op::Report(Quantum::Completed),
Op::Service,
]);
assert_eq!(model.registered(), 0);
model.run(&[Op::Detach(0)]);
}
#[test]
fn claiming_is_single_flight() {
let mut model = Model::new(8, 8);
model.run(&[Op::Begin, Op::Service, Op::Claim]);
assert_eq!(model.claimed(), Some(model.id(0)));
model.run(&[Op::Claim]);
assert_eq!(model.claimed(), None);
model.run(&[Op::Report(Quantum::Idle), Op::Schedule(0), Op::Claim]);
assert_eq!(model.claimed(), Some(model.id(0)));
model.run(&[Op::Report(Quantum::Idle), Op::Detach(0)]);
}
#[test]
fn scheduling_is_idempotent() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Schedule(0),
Op::Schedule(0),
Op::Schedule(0),
]);
assert_eq!(model.ready(), 1, "one entry, however often it is scheduled");
model.run(&[Op::Claim]);
assert_eq!(model.ready(), 0);
model.run(&[Op::Claim]);
assert_eq!(model.claimed(), None, "the queue held only the one");
model.run(&[Op::Detach(0)]);
}
#[test]
fn scheduling_a_claimed_enumeration_does_not_queue_it() {
let mut model = Model::new(8, 8);
model.run(&[Op::Begin, Op::Service, Op::Claim, Op::Schedule(0)]);
assert_eq!(
model.ready(),
0,
"re-queuing it would let a second worker take the same buffer"
);
model.run(&[Op::Report(Quantum::Idle), Op::Detach(0)]);
}
#[test]
fn a_finished_quantum_outranks_a_concurrent_cancellation() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::Cancel(0),
Op::Service,
Op::Report(Quantum::Completed),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.terminal(0), Some("completed"));
assert_eq!(model.registered(), 0);
}
#[test]
fn a_failed_quantum_delivers_a_failed_terminal() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::OfferEntry(0, "before"),
Op::RunEngine(Quantum::Failed),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(
model.entries(0),
["before"],
"a late failure truncates rather than retracts"
);
assert_eq!(model.terminal(0), Some("failed"));
assert_eq!(model.registered(), 0);
model.run(&[Op::Detach(0)]);
}
#[test]
fn a_worker_may_report_cancellation_as_its_own_outcome() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::RunEngine(Quantum::Cancelled),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.terminal(0), Some("cancelled"));
assert_eq!(model.registered(), 0);
model.run(&[Op::Detach(0)]);
}
#[test]
fn a_parked_quantum_is_re_queued_when_room_appears() {
let mut model = Model::new(8, 3);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::OfferEntry(0, "a"),
Op::OfferEntry(0, "b"),
Op::Report(Quantum::Parked),
]);
assert_eq!(model.ready(), 0, "parked, not runnable");
model.run(&[Op::Recv]);
assert_eq!(model.ready(), 1, "taking a record made it runnable again");
model.run(&[
Op::Claim,
Op::OfferEntry(0, "c"),
Op::Report(Quantum::Completed),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.entries(0), ["a", "b", "c"]);
assert_eq!(model.terminal(0), Some("completed"));
model.run(&[Op::Detach(0)]);
}
#[test]
fn servicing_a_begin_makes_it_runnable() {
let mut model = Model::new(8, 8);
model.run(&[Op::Begin]);
assert_eq!(model.ready(), 0, "not registered yet");
model.run(&[Op::Service]);
assert_eq!(model.ready(), 1);
model.run(&[Op::Detach(0)]);
}
#[test]
fn the_minimum_submission_ring_covers_cancel_and_retire() {
let mut model = Model::new(MINIMUM_SUBMISSION_CAPACITY, 8);
model.run(&[
Op::Begin,
Op::BeginRefused(BeginFailure::SubmissionRingFull),
Op::Service,
Op::BeginRefused(BeginFailure::SubmissionRingFull),
Op::RunEngine(Quantum::Completed),
Op::Service,
Op::Detach(0),
Op::Begin,
Op::Service,
Op::Detach(1),
]);
assert_eq!(model.registered(), 1);
}
#[test]
fn a_yielding_quantum_is_re_queued() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::Report(Quantum::Yielded),
]);
assert_eq!(model.ready(), 1, "yielding asks for another turn");
model.run(&[Op::Claim, Op::Report(Quantum::Idle)]);
assert_eq!(model.ready(), 0, "an idle quantum does not");
model.run(&[Op::Detach(0)]);
}
#[test]
fn a_cancelled_enumeration_does_not_yield_again() {
let mut model = Model::new(8, 8);
model.run(&[
Op::Begin,
Op::Service,
Op::Claim,
Op::Cancel(0),
Op::Service,
Op::Report(Quantum::Yielded),
Op::Service,
Op::DrainReceiver,
]);
assert_eq!(model.terminal(0), Some("cancelled"));
assert_eq!(model.ready(), 0);
assert_eq!(model.registered(), 0);
}