use async_nats::HeaderMap;
use routers_network::Entry;
use tokio::time::Instant;
use web_time::SystemTime;
use crate::bus::Wire;
use crate::bus::adapter::{AckHandle, Delivery};
use crate::event::Payload;
use crate::ingress;
use crate::orchestrator::frontier::{FrontierTracker, Observed};
use crate::orchestrator::scheduler::{
Enqueue, EnqueueReadiness, PendingObservation, PublishedAtMicros, Scheduler,
};
use crate::partition;
use crate::protocol::ids::{ObservationId, SCHEMA_VERSION, SchemaVersion, headers};
use crate::topology::partition_of_subject;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum PoisonReason {
Decode,
SubjectPartition {
expected: u16,
got: Option<u16>,
},
VehiclePartition {
subject: u16,
vehicle: u16,
},
Schema {
got: Option<u32>,
},
Invalid {
kind: &'static str,
},
PublishTime,
}
impl PoisonReason {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
PoisonReason::Decode => "decode",
PoisonReason::SubjectPartition { .. } => "subject_partition",
PoisonReason::VehiclePartition { .. } => "vehicle_partition",
PoisonReason::Schema { .. } => "schema",
PoisonReason::Invalid { kind } => kind,
PoisonReason::PublishTime => "publish_time",
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum SuppressReason {
BehindFrontier,
Committed,
}
impl SuppressReason {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
SuppressReason::BehindFrontier => "behind_frontier",
SuppressReason::Committed => "committed",
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum DeferReason {
PendingFull,
DuplicateOwned,
}
impl DeferReason {
#[must_use]
pub fn label(self) -> &'static str {
match self {
DeferReason::PendingFull => "pending_full",
DeferReason::DuplicateOwned => "duplicate_owned",
}
}
}
#[derive(Debug)]
#[must_use = "a disposition may carry an ack handle that must be acknowledged"]
pub enum RawDisposition<H: AckHandle> {
Queued {
depth: usize,
},
Poison {
handle: H,
reason: PoisonReason,
},
Suppressed {
handle: H,
reason: SuppressReason,
},
Deferred {
handle: H,
reason: DeferReason,
},
}
impl<H: AckHandle> RawDisposition<H> {
#[must_use]
pub fn is_terminal(&self) -> bool {
match self {
RawDisposition::Poison { .. } => true,
RawDisposition::Suppressed { .. } => true,
RawDisposition::Queued { .. } => false,
RawDisposition::Deferred { .. } => false,
}
}
}
pub struct RawEnvelope<'a, H: AckHandle> {
pub subject: &'a str,
pub headers: Option<&'a HeaderMap>,
pub bytes: &'a [u8],
pub sent_at: Option<SystemTime>,
pub handle: H,
}
struct DecodedRaw<H: AckHandle> {
schema: Option<SchemaVersion>,
payload: Payload,
published_at: Option<SystemTime>,
handle: H,
}
#[derive(Clone, Copy, Debug)]
pub struct RawReader {
partition: u16,
}
impl RawReader {
#[must_use]
pub fn new(partition: u16) -> Self {
Self { partition }
}
#[must_use]
pub fn partition(&self) -> u16 {
self.partition
}
pub fn admit<E, H>(
&self,
scheduler: &mut Scheduler<E, H>,
tracker: &mut FrontierTracker,
delivery: Delivery<Payload, H>,
now: Instant,
) -> RawDisposition<H>
where
E: Entry,
H: AckHandle,
{
let observed = tracker.observe(delivery.handle.sequence());
let got = partition_of_subject(&delivery.subject);
if got != Some(self.partition) {
return RawDisposition::Poison {
handle: delivery.handle,
reason: PoisonReason::SubjectPartition {
expected: self.partition,
got,
},
};
}
self.admit_payload(
scheduler,
tracker,
DecodedRaw {
schema: None,
payload: delivery.item,
published_at: delivery.sent_at,
handle: delivery.handle,
},
now,
observed,
)
}
pub fn admit_bytes<E, H>(
&self,
scheduler: &mut Scheduler<E, H>,
tracker: &mut FrontierTracker,
envelope: RawEnvelope<'_, H>,
now: Instant,
) -> RawDisposition<H>
where
E: Entry,
H: AckHandle,
{
let RawEnvelope {
subject,
headers,
bytes,
sent_at,
handle,
} = envelope;
let observed = tracker.observe(handle.sequence());
let got = partition_of_subject(subject);
if got != Some(self.partition) {
return RawDisposition::Poison {
handle,
reason: PoisonReason::SubjectPartition {
expected: self.partition,
got,
},
};
}
let payload = match Payload::decode(bytes) {
Ok(payload) => payload,
Err(_) => {
return RawDisposition::Poison {
handle,
reason: PoisonReason::Decode,
};
}
};
self.admit_payload(
scheduler,
tracker,
DecodedRaw {
schema: headers.and_then(headers::schema_of),
payload,
published_at: sent_at,
handle,
},
now,
observed,
)
}
fn admit_payload<E, H>(
&self,
scheduler: &mut Scheduler<E, H>,
tracker: &mut FrontierTracker,
raw: DecodedRaw<H>,
now: Instant,
observed: Observed,
) -> RawDisposition<H>
where
E: Entry,
H: AckHandle,
{
let DecodedRaw {
schema,
payload,
published_at,
handle,
} = raw;
let vehicle_partition = partition::partition_of(payload.vehicle_id) as u16;
if vehicle_partition != self.partition {
return RawDisposition::Poison {
handle,
reason: PoisonReason::VehiclePartition {
subject: self.partition,
vehicle: vehicle_partition,
},
};
}
if let Some(version) = schema
&& version != SCHEMA_VERSION
{
return RawDisposition::Poison {
handle,
reason: PoisonReason::Schema {
got: Some(version.0),
},
};
}
if let Err(error) = ingress::validate(&payload) {
return RawDisposition::Poison {
handle,
reason: PoisonReason::Invalid { kind: error.kind() },
};
}
let Some(published_at) = published_at.and_then(PublishedAtMicros::from_system_time) else {
return RawDisposition::Poison {
handle,
reason: PoisonReason::PublishTime,
};
};
let sequence = handle.sequence();
if sequence <= tracker.frontier() {
return RawDisposition::Suppressed {
handle,
reason: SuppressReason::BehindFrontier,
};
}
let observation = PendingObservation {
id: ObservationId {
partition: self.partition,
sequence,
},
payload,
published_at,
handle,
received: now,
};
let readiness = scheduler.enqueue_readiness(&observation);
match observed {
Observed::BehindFrontier => {
return RawDisposition::Suppressed {
handle: observation.handle,
reason: SuppressReason::BehindFrontier,
};
}
Observed::Duplicate if readiness == EnqueueReadiness::Coalesced => {
return RawDisposition::Deferred {
handle: observation.handle,
reason: DeferReason::DuplicateOwned,
};
}
Observed::Duplicate | Observed::New => {}
}
if matches!(readiness, EnqueueReadiness::Backpressured { .. }) {
return RawDisposition::Deferred {
handle: observation.handle,
reason: DeferReason::PendingFull,
};
}
match scheduler.enqueue(observation, now) {
Enqueue::Queued { depth } => RawDisposition::Queued { depth },
Enqueue::Coalesced(handle) => RawDisposition::Deferred {
handle,
reason: DeferReason::DuplicateOwned,
},
Enqueue::Committed(handle) => RawDisposition::Suppressed {
handle,
reason: SuppressReason::Committed,
},
Enqueue::Overflow { handle, .. } => RawDisposition::Deferred {
handle,
reason: DeferReason::PendingFull,
},
}
}
}
#[cfg(test)]
mod tests {
use alloc::sync::Arc;
use core::time::Duration;
use std::sync::Mutex;
use chrono::{TimeZone, Utc};
use geo::Point;
use routers_network::mock::MockEntryId;
use super::*;
use crate::event::VehicleId;
use crate::orchestrator::scheduler::{CheckpointState, SchedulerConfig};
use crate::protocol::ids::{GraphVersion, RegionId, Revision, SCHEMA_VERSION, SegmentId};
use crate::store::checkpoint::{PartitionFrontier, VehicleCheckpoint};
use web_time::UNIX_EPOCH;
use routers_transition::matcher::Trip;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum AckOp {
Ack(u64),
Nak(u64),
}
#[derive(Clone, Debug)]
struct TestAck {
seq: u64,
log: Arc<Mutex<Vec<AckOp>>>,
}
impl AckHandle for TestAck {
async fn ack(self) -> anyhow::Result<()> {
self.log.lock().unwrap().push(AckOp::Ack(self.seq));
Ok(())
}
async fn nak(self, _delay: Option<Duration>) -> anyhow::Result<()> {
self.log.lock().unwrap().push(AckOp::Nak(self.seq));
Ok(())
}
fn sequence(&self) -> u64 {
self.seq
}
fn deliveries(&self) -> u32 {
1
}
}
const PARTITION: u16 = 485;
const VEHICLE: u64 = 1;
struct Fixture {
log: Arc<Mutex<Vec<AckOp>>>,
region: RegionId,
now: Instant,
}
impl Fixture {
fn new() -> Self {
Self {
log: Arc::new(Mutex::new(Vec::new())),
region: RegionId::new("r1").unwrap(),
now: Instant::now(),
}
}
fn ack(&self, seq: u64) -> TestAck {
TestAck {
seq,
log: self.log.clone(),
}
}
fn payload(&self, vehicle: u64, lon: f64, lat: f64) -> Payload {
Payload {
vehicle_id: VehicleId(vehicle),
timestamp: Utc.timestamp_micros(1_775_000_000_000_000).unwrap(),
point: Point::new(lon, lat),
}
}
fn good(&self) -> Payload {
self.payload(VEHICLE, 151.2093, -33.8688)
}
fn envelope<'a>(
&self,
subject: &'a str,
headers: Option<&'a HeaderMap>,
bytes: &'a [u8],
seq: u64,
) -> RawEnvelope<'a, TestAck> {
RawEnvelope {
subject,
headers,
bytes,
sent_at: Some(UNIX_EPOCH + Duration::from_secs(1)),
handle: self.ack(seq),
}
}
fn present_checkpoint(&self, last_seq: u64) -> CheckpointState<MockEntryId> {
CheckpointState::Present(VehicleCheckpoint {
trip: Trip::new(),
last_input: ObservationId {
partition: PARTITION,
sequence: last_seq,
},
revision: Revision(last_seq),
segment: SegmentId(1),
finalized_through: None,
graph: GraphVersion::new("g1").unwrap(),
schema: SCHEMA_VERSION,
region: self.region.clone(),
routing_version: 1,
})
}
}
fn scheduler() -> Scheduler<MockEntryId, TestAck> {
Scheduler::new(SchedulerConfig::default())
}
fn tracker(fx: &Fixture) -> FrontierTracker {
FrontierTracker::new(PARTITION, None, fx.now)
}
fn subject() -> String {
crate::topology::raw_subject(u64::from(PARTITION))
}
#[test]
fn happy_path_queues_and_observes() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 10),
fx.now,
);
assert!(matches!(disposition, RawDisposition::Queued { depth: 1 }));
assert!(!disposition.is_terminal());
assert_eq!(sched.stats().pending, 1);
assert_eq!(frontier.outstanding(), 1);
assert_eq!(frontier.oldest_outstanding(), Some(10));
}
#[test]
fn decoded_delivery_path_queues() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let delivery = Delivery {
item: fx.good(),
handle: fx.ack(4),
subject: subject(),
msg_id: None,
headers: HeaderMap::new(),
sent_at: Some(UNIX_EPOCH + Duration::from_secs(1)),
redelivered: false,
};
let disposition = reader.admit(&mut sched, &mut frontier, delivery, fx.now);
assert!(matches!(disposition, RawDisposition::Queued { depth: 1 }));
assert_eq!(sched.stats().pending, 1);
}
#[test]
fn wrong_subject_partition_is_poison() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = crate::topology::raw_subject(486);
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 1),
fx.now,
);
match disposition {
RawDisposition::Poison {
handle,
reason: PoisonReason::SubjectPartition { expected, got },
} => {
assert_eq!(handle.sequence(), 1);
assert_eq!(expected, PARTITION);
assert_eq!(got, Some(486));
}
other => panic!("expected SubjectPartition poison, got {other:?}"),
}
assert_eq!(
frontier.outstanding(),
1,
"even poison reserves a frontier hole until the worker acks it"
);
assert_eq!(sched.stats().pending, 0);
}
#[test]
fn missing_publish_time_is_poison_after_reserving_the_frontier_hole() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let mut envelope = fx.envelope(&subject, None, &bytes, 3);
envelope.sent_at = None;
let disposition = reader.admit_bytes(&mut sched, &mut frontier, envelope, fx.now);
assert!(matches!(
disposition,
RawDisposition::Poison {
reason: PoisonReason::PublishTime,
..
}
));
assert_eq!(frontier.oldest_outstanding(), Some(3));
assert_eq!(sched.stats().pending, 0);
}
#[test]
fn unparseable_subject_is_poison() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope("events.raw.p.not-a-number", None, &bytes, 1),
fx.now,
);
assert!(matches!(
disposition,
RawDisposition::Poison {
reason: PoisonReason::SubjectPartition { got: None, .. },
..
}
));
}
#[test]
fn undecodable_payload_is_poison() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let subject = subject();
let bytes = [0x0a, 0x7f];
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 2),
fx.now,
);
assert!(disposition.is_terminal());
match disposition {
RawDisposition::Poison {
handle,
reason: PoisonReason::Decode,
} => assert_eq!(handle.sequence(), 2),
other => panic!("expected Decode poison, got {other:?}"),
}
}
#[test]
fn wrong_vehicle_partition_is_poison() {
let fx = Fixture::new();
let reader = RawReader::new(0);
let mut sched = scheduler();
let mut frontier = FrontierTracker::new(0, None, fx.now);
let bytes = fx.good().encode().unwrap();
let subject = crate::topology::raw_subject(0);
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 1),
fx.now,
);
match disposition {
RawDisposition::Poison {
reason: PoisonReason::VehiclePartition { subject, vehicle },
..
} => {
assert_eq!(subject, 0);
assert_eq!(vehicle, PARTITION);
}
other => panic!("expected VehiclePartition poison, got {other:?}"),
}
}
#[test]
fn schema_mismatch_is_poison() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let mut wrong = HeaderMap::new();
wrong.insert(headers::SCHEMA, (SCHEMA_VERSION.0 + 1).to_string());
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, Some(&wrong), &bytes, 1),
fx.now,
);
match disposition {
RawDisposition::Poison {
reason: PoisonReason::Schema { got },
..
} => assert_eq!(got, Some(SCHEMA_VERSION.0 + 1)),
other => panic!("expected Schema poison, got {other:?}"),
}
}
#[test]
fn matching_schema_header_passes() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let mut ok = HeaderMap::new();
headers::stamp_schema(&mut ok);
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, Some(&ok), &bytes, 1),
fx.now,
);
assert!(matches!(disposition, RawDisposition::Queued { .. }));
}
#[test]
fn zero_vehicle_is_poison_invalid() {
let fx = Fixture::new();
let reader = RawReader::new(0);
let mut sched = scheduler();
let mut frontier = FrontierTracker::new(0, None, fx.now);
let bytes = fx.payload(0, 151.0, -33.0).encode().unwrap();
let subject = crate::topology::raw_subject(0);
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 1),
fx.now,
);
match disposition {
RawDisposition::Poison {
reason: PoisonReason::Invalid { kind },
..
} => assert_eq!(kind, "zero_vehicle"),
other => panic!("expected Invalid poison, got {other:?}"),
}
}
#[test]
fn non_finite_coordinate_is_poison_invalid() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.payload(VEHICLE, f64::NAN, -33.0).encode().unwrap();
let subject = subject();
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 1),
fx.now,
);
match disposition {
RawDisposition::Poison {
reason: PoisonReason::Invalid { kind },
..
} => assert_eq!(kind, "non_finite"),
other => panic!("expected Invalid poison, got {other:?}"),
}
}
#[test]
fn behind_frontier_is_suppressed_and_terminal() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = FrontierTracker::new(
PARTITION,
Some(PartitionFrontier {
partition: PARTITION,
sequence: 100,
}),
fx.now,
);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 50),
fx.now,
);
assert!(disposition_terminal(&disposition));
match disposition {
RawDisposition::Suppressed {
handle,
reason: SuppressReason::BehindFrontier,
} => assert_eq!(handle.sequence(), 50),
other => panic!("expected BehindFrontier suppression, got {other:?}"),
}
assert_eq!(frontier.outstanding(), 0);
assert_eq!(sched.stats().pending, 0);
}
#[test]
fn duplicate_in_flight_is_coalesced_and_not_terminal() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let first = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 7),
fx.now,
);
assert!(matches!(first, RawDisposition::Queued { .. }));
let again = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 7),
fx.now,
);
match &again {
RawDisposition::Deferred {
handle,
reason: DeferReason::DuplicateOwned,
} => assert_eq!(handle.sequence(), 7),
other => panic!("expected DuplicateOwned deferral, got {other:?}"),
}
assert!(!again.is_terminal());
assert_eq!(sched.stats().pending, 1);
assert_eq!(frontier.outstanding(), 1);
}
#[test]
fn already_committed_is_suppressed_and_terminal() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let seed = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 10),
fx.now,
);
assert!(matches!(seed, RawDisposition::Queued { .. }));
sched.set_checkpoint(VehicleId(VEHICLE), fx.present_checkpoint(100));
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 50),
fx.now,
);
match &disposition {
RawDisposition::Suppressed {
handle,
reason: SuppressReason::Committed,
} => assert_eq!(handle.sequence(), 50),
other => panic!("expected Committed suppression, got {other:?}"),
}
assert!(disposition.is_terminal());
}
#[test]
fn full_pending_queue_defers_and_reserves_the_frontier_hole() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = Scheduler::<MockEntryId, TestAck>::new(SchedulerConfig {
pending_limit: 1,
..SchedulerConfig::default()
});
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = subject();
let first = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 10),
fx.now,
);
assert!(matches!(first, RawDisposition::Queued { depth: 1 }));
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 11),
fx.now,
);
match &disposition {
RawDisposition::Deferred {
handle,
reason: DeferReason::PendingFull,
} => assert_eq!(handle.sequence(), 11),
other => panic!("expected PendingFull deferral, got {other:?}"),
}
assert!(!disposition.is_terminal());
assert_eq!(
frontier.outstanding(),
2,
"the deferred raw reserves a hole"
);
assert_eq!(frontier.oldest_outstanding(), Some(10));
}
#[test]
fn deferred_hole_blocks_later_completion_and_redelivery_can_enqueue() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut full = Scheduler::<MockEntryId, TestAck>::new(SchedulerConfig {
pending_limit: 1,
..SchedulerConfig::default()
});
let mut open = scheduler();
let mut frontier = FrontierTracker::new(
PARTITION,
Some(PartitionFrontier {
partition: PARTITION,
sequence: 10,
}),
fx.now,
);
let subject = subject();
let first_payload = fx.good();
let first_bytes = first_payload.encode().unwrap();
assert!(matches!(
full.enqueue(
PendingObservation {
id: ObservationId {
partition: PARTITION,
sequence: 9,
},
payload: first_payload,
published_at: PublishedAtMicros::from_unix_micros(1_000_000)
.expect("test timestamp is non-negative"),
handle: fx.ack(9),
received: fx.now,
},
fx.now,
),
Enqueue::Queued { .. }
));
let deferred = reader.admit_bytes(
&mut full,
&mut frontier,
fx.envelope(&subject, None, &first_bytes, 11),
fx.now,
);
assert!(matches!(
deferred,
RawDisposition::Deferred {
reason: DeferReason::PendingFull,
..
}
));
let other_vehicle = (VEHICLE + 1..)
.find(|&vehicle| partition::partition_of(VehicleId(vehicle)) as u16 == PARTITION)
.expect("another vehicle in the partition");
let later_bytes = fx.payload(other_vehicle, 151.21, -33.86).encode().unwrap();
assert!(matches!(
reader.admit_bytes(
&mut open,
&mut frontier,
fx.envelope(&subject, None, &later_bytes, 12),
fx.now,
),
RawDisposition::Queued { .. }
));
frontier.complete(12);
assert_eq!(frontier.frontier(), 10, "sequence 11 remains a hard hole");
assert!(matches!(
reader.admit_bytes(
&mut open,
&mut frontier,
fx.envelope(&subject, None, &first_bytes, 11),
fx.now,
),
RawDisposition::Queued { .. }
));
frontier.complete(11);
assert_eq!(frontier.frontier(), 12);
}
#[test]
fn returned_handle_can_be_acknowledged() {
let fx = Fixture::new();
let reader = RawReader::new(PARTITION);
let mut sched = scheduler();
let mut frontier = tracker(&fx);
let bytes = fx.good().encode().unwrap();
let subject = crate::topology::raw_subject(486);
let disposition = reader.admit_bytes(
&mut sched,
&mut frontier,
fx.envelope(&subject, None, &bytes, 9),
fx.now,
);
let RawDisposition::Poison { handle, .. } = disposition else {
panic!("expected poison");
};
futures::executor::block_on(handle.ack()).unwrap();
assert_eq!(&*fx.log.lock().unwrap(), &[AckOp::Ack(9)]);
}
#[test]
fn labels_are_bounded() {
assert_eq!(PoisonReason::Decode.label(), "decode");
assert_eq!(
PoisonReason::SubjectPartition {
expected: 0,
got: None
}
.label(),
"subject_partition"
);
assert_eq!(
PoisonReason::VehiclePartition {
subject: 0,
vehicle: 1
}
.label(),
"vehicle_partition"
);
assert_eq!(PoisonReason::Schema { got: None }.label(), "schema");
assert_eq!(PoisonReason::PublishTime.label(), "publish_time");
assert_eq!(
PoisonReason::Invalid {
kind: "zero_vehicle"
}
.label(),
"zero_vehicle"
);
assert_eq!(SuppressReason::BehindFrontier.label(), "behind_frontier");
assert_eq!(SuppressReason::Committed.label(), "committed");
assert_eq!(DeferReason::PendingFull.label(), "pending_full");
assert_eq!(DeferReason::DuplicateOwned.label(), "duplicate_owned");
}
fn disposition_terminal(d: &RawDisposition<TestAck>) -> bool {
d.is_terminal()
}
}