use alloc::collections::VecDeque;
use alloc::sync::Arc;
use core::time::Duration;
use std::collections::{HashMap, HashSet};
use futures::future::BoxFuture;
use futures::stream::{FuturesUnordered, StreamExt as _};
use routers_network::Entry;
use tokio::time::{Instant, MissedTickBehavior, interval};
use tracing::{debug, error, warn};
use web_time::UNIX_EPOCH;
use crate::bus::adapter::{AckHandle, Delivery, Publisher, Source};
use crate::event::VehicleId;
use crate::lifecycle::{Drain, Shutdown};
use crate::matcher::pull::RawBytes;
use crate::metrics::Metrics;
use crate::orchestrator::admission::{Admission, HeldReason, Waiting};
use crate::orchestrator::commit::{
self, CommitConfig, CommitError, CommitMode, Committed, Committer, Decision,
};
use crate::orchestrator::dispatch::{
DispatchConfig, DispatchError, Dispatcher, RetryRequest, STICKY_HEADER, TimestampRegression,
timestamp_regression,
};
use crate::orchestrator::frontier::{FrontierConfig, FrontierTracker};
use crate::orchestrator::reader::{DeferReason, RawDisposition, RawEnvelope, RawReader};
use crate::orchestrator::recovery::{RecoveryReport, restore_vehicle};
use crate::orchestrator::scheduler::{
ActiveJob, CheckpointState, JobReservation, Scheduler, SchedulerConfig, VehicleState,
};
use crate::orchestrator::validate::{self, QuarantineReason, RejectReason, Verdict};
use crate::protocol::ids::{GraphVersion, JobId, RegionId, Revision, SCHEMA_VERSION, SegmentId};
use crate::protocol::job::{BaseState, JobIdentity, SolveRequest};
use crate::protocol::output::{CommittedOutput, ResetReason, TerminalReason};
use crate::protocol::result::{SolveOutcome, SolveResult};
use crate::region::catalog::Catalog;
use crate::region::resolver::{Pin, Resolver};
use crate::store::checkpoint::{
CheckpointStore, DirectOutcome, PartitionFrontier, VehicleCheckpoint,
};
#[derive(Clone, Debug)]
pub struct WorkerConfig {
pub partition: u16,
pub scheduler: SchedulerConfig,
pub parked_limit: usize,
pub frontier: FrontierConfig,
pub dispatch: DispatchConfig,
pub commit: CommitConfig,
pub tick: Duration,
pub evict_every: Duration,
pub blocked_retry: Duration,
pub request_retry: Duration,
pub sticky_retry: Duration,
pub grace: Duration,
pub checkpoint: CheckpointPolicy,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct CheckpointPolicy {
pub every: usize,
pub interval: Duration,
pub in_flight: usize,
}
impl Default for CheckpointPolicy {
fn default() -> Self {
Self {
every: 16,
interval: Duration::from_secs(2),
in_flight: 32,
}
}
}
impl WorkerConfig {
#[must_use]
pub fn new(partition: u16) -> Self {
Self {
partition,
scheduler: SchedulerConfig::default(),
parked_limit: 4,
frontier: FrontierConfig::default(),
dispatch: DispatchConfig::default(),
commit: CommitConfig::default(),
tick: Duration::from_millis(250),
evict_every: Duration::from_secs(30),
blocked_retry: Duration::from_secs(5),
request_retry: Duration::from_secs(3),
sticky_retry: Duration::from_millis(500),
checkpoint: CheckpointPolicy::default(),
grace: Duration::from_secs(20),
}
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct WorkerStats {
pub observed: u64,
pub queued: u64,
pub coalesced: u64,
pub suppressed: u64,
pub deferred: u64,
pub poison: u64,
pub dispatched: u64,
pub held: u64,
pub accepted: u64,
pub parked: u64,
pub rejected: u64,
pub quarantined: u64,
pub quarantined_vehicles: u64,
pub committed: u64,
pub terminal: u64,
pub resets: u64,
pub conflicts: u64,
pub frontier: u64,
pub persisted: u64,
}
#[derive(Clone)]
struct ActiveMeta {
reset: Option<ResetReason>,
segment: SegmentId,
region: RegionId,
graph: GraphVersion,
routing_version: u64,
expected_base: Option<Revision>,
}
#[derive(Clone, Copy)]
struct PendingReset {
reason: ResetReason,
prior: Option<Revision>,
}
struct LocalTerminalRoute {
region: RegionId,
graph: GraphVersion,
routing_version: u64,
}
struct ParkedResult<E: Entry, H: AckHandle> {
result: SolveResult<E>,
handle: H,
}
struct CommitCompletion<E: Entry, H: AckHandle, SE> {
vehicle: VehicleId,
next: VehicleCheckpoint<E>,
result_delivery: Option<ParkedResult<E, H>>,
started: Instant,
completed: Instant,
sampled: bool,
is_terminal: bool,
terminal_reason: Option<&'static str>,
reset_reason: Option<&'static str>,
result: Result<Committed, CommitError<SE>>,
}
type CommitFuture<E, H, SE> = BoxFuture<'static, CommitCompletion<E, H, SE>>;
struct Dirty {
sequences: Vec<u64>,
since: Instant,
in_flight: bool,
retry_at: Option<Instant>,
}
struct PersistCompletion<SE> {
vehicle: VehicleId,
revision: Revision,
sequences: Vec<u64>,
started: Instant,
result: Result<DirectOutcome, CommitError<SE>>,
}
type PersistFuture<SE> = BoxFuture<'static, PersistCompletion<SE>>;
enum ResultVerdict {
Accept,
Park,
Reject(RejectReason),
Quarantine(QuarantineReason),
}
impl From<Verdict<'_>> for ResultVerdict {
fn from(verdict: Verdict<'_>) -> Self {
match verdict {
Verdict::Accept { .. } => ResultVerdict::Accept,
Verdict::Park => ResultVerdict::Park,
Verdict::Reject(reason) => ResultVerdict::Reject(reason),
Verdict::Quarantine(reason) => ResultVerdict::Quarantine(reason),
}
}
}
enum DispatchOutcome {
Done,
RegionHeld,
GloballyHeld,
Skipped,
}
pub struct PartitionWorker<E, S, JP, OP, RS, XS>
where
E: Entry + serde::de::DeserializeOwned,
S: CheckpointStore,
RS: Source<RawBytes>,
XS: Source<SolveResult<E>>,
{
cfg: WorkerConfig,
catalog: Arc<Catalog>,
admission: Admission,
store: S,
dispatcher: Dispatcher<JP>,
committer: Committer<E, S, OP>,
reader: RawReader,
raw: RS,
results: XS,
scheduler: Scheduler<E, RS::Handle>,
tracker: FrontierTracker,
blocked: HashMap<VehicleId, Instant>,
held: HashMap<VehicleId, Waiting>,
active_meta: HashMap<VehicleId, ActiveMeta>,
requests: HashMap<VehicleId, (RetryRequest<E>, Instant)>,
sticky: HashMap<VehicleId, (String, Revision)>,
pending_reset: HashMap<VehicleId, PendingReset>,
parked_results: HashMap<VehicleId, Vec<ParkedResult<E, XS::Handle>>>,
blocked_results: HashMap<VehicleId, ParkedResult<E, XS::Handle>>,
quarantined: HashSet<VehicleId>,
commits: FuturesUnordered<CommitFuture<E, XS::Handle, S::Error>>,
dirty: HashMap<VehicleId, Dirty>,
unpersisted: HashSet<u64>,
persisted_base: HashMap<VehicleId, Option<Revision>>,
persists: FuturesUnordered<PersistFuture<S::Error>>,
epoch: Option<u64>,
fenced: Option<u64>,
last_heartbeat: Instant,
shutdown: Shutdown,
drain: Drain,
last_evict: Instant,
stats: WorkerStats,
metrics: Metrics,
}
impl<E, S, JP, OP, RS, XS> PartitionWorker<E, S, JP, OP, RS, XS>
where
E: Entry + serde::de::DeserializeOwned,
S: CheckpointStore,
JP: Publisher<SolveRequest<E>>,
OP: Publisher<CommittedOutput<E>>,
RS: Source<RawBytes>,
XS: Source<SolveResult<E>>,
{
#[allow(clippy::too_many_arguments)]
#[must_use]
pub fn new(
cfg: WorkerConfig,
catalog: Arc<Catalog>,
admission: Admission,
store: S,
dispatcher: Dispatcher<JP>,
committer: Committer<E, S, OP>,
raw: RS,
results: XS,
report: &RecoveryReport,
shutdown: Shutdown,
) -> Self {
let now = Instant::now();
let initial = report.frontier.map(|sequence| PartitionFrontier {
partition: cfg.partition,
sequence,
});
let tracker = FrontierTracker::with_config(cfg.partition, initial, now, cfg.frontier);
let scheduler = Scheduler::new(cfg.scheduler);
let reader = RawReader::new(cfg.partition);
let mut blocked = HashMap::new();
for &vehicle in &report.prepared_failed {
blocked.insert(vehicle, now);
}
Self {
cfg,
catalog,
admission,
store,
dispatcher,
committer: committer.with_epoch(report.epoch),
reader,
raw,
results,
scheduler,
tracker,
blocked,
held: HashMap::new(),
active_meta: HashMap::new(),
requests: HashMap::new(),
sticky: HashMap::new(),
pending_reset: HashMap::new(),
parked_results: HashMap::new(),
blocked_results: HashMap::new(),
quarantined: HashSet::new(),
commits: FuturesUnordered::new(),
dirty: HashMap::new(),
unpersisted: HashSet::new(),
persisted_base: HashMap::new(),
persists: FuturesUnordered::new(),
epoch: report.epoch,
fenced: None,
last_heartbeat: now,
shutdown,
drain: Drain::new(),
last_evict: now,
stats: WorkerStats::default(),
metrics: Metrics::noop(),
}
}
#[must_use]
pub fn with_metrics(mut self, metrics: Metrics) -> Self {
self.committer = self.committer.with_metrics(metrics.clone());
self.metrics = metrics;
self
}
pub async fn run(mut self) -> anyhow::Result<WorkerStats> {
let mut tick = interval(self.cfg.tick);
tick.set_missed_tick_behavior(MissedTickBehavior::Delay);
let mut raw_open = true;
let mut results_open = true;
loop {
if let Some(owner) = self.fenced {
anyhow::bail!(
"partition {} ownership lost to epoch {owner}",
self.cfg.partition
);
}
tokio::select! {
() = self.shutdown.triggered() => break,
maybe = self.raw.next(), if raw_open => {
let started = Instant::now();
match maybe {
Some(Ok(delivery)) => self.on_raw(delivery).await,
Some(Err(error)) => warn!(%error, "raw source error"),
None => raw_open = false,
}
self.metrics.worker_turn_seconds("raw", started.elapsed().as_secs_f64());
},
maybe = self.results.next(), if results_open => {
let started = Instant::now();
match maybe {
Some(Ok(delivery)) => self.on_result(delivery).await,
Some(Err(error)) => warn!(%error, "result source error"),
None => results_open = false,
}
self.metrics.worker_turn_seconds("result", started.elapsed().as_secs_f64());
},
completion = self.commits.next(), if !self.commits.is_empty() => {
let started = Instant::now();
self.finish_commit(completion.expect("guarded by non-empty check")).await;
self.pump().await;
self.metrics.worker_turn_seconds("commit", started.elapsed().as_secs_f64());
},
completion = self.persists.next(), if !self.persists.is_empty() => {
let started = Instant::now();
self.finish_persist(completion.expect("guarded by non-empty check")).await;
self.metrics.worker_turn_seconds("persist", started.elapsed().as_secs_f64());
},
_ = tick.tick() => {
let started = Instant::now();
self.on_tick().await;
self.metrics.worker_turn_seconds("tick", started.elapsed().as_secs_f64());
},
}
}
let deadline = Instant::now() + self.cfg.grace;
while !self.commits.is_empty() {
let remaining = deadline.saturating_duration_since(Instant::now());
let Ok(Some(completion)) = tokio::time::timeout(remaining, self.commits.next()).await
else {
break;
};
self.finish_commit(completion).await;
}
let vehicles: Vec<VehicleId> = self.dirty.keys().copied().collect();
for vehicle in vehicles {
self.start_persist(vehicle, Instant::now(), true);
}
while !self.persists.is_empty() {
let remaining = deadline.saturating_duration_since(Instant::now());
let Ok(Some(completion)) = tokio::time::timeout(remaining, self.persists.next()).await
else {
break;
};
self.finish_persist(completion).await;
}
let _ = self.drain.quiesce(Duration::ZERO).await;
let frontier = PartitionFrontier {
partition: self.cfg.partition,
sequence: self.tracker.frontier(),
};
if let Ok(false) = self.write_frontier(frontier).await {
self.fenced = self.fenced.or(Some(0));
}
self.stats.frontier = self.tracker.frontier();
if let Some(owner) = self.fenced {
anyhow::bail!(
"partition {} ownership lost to epoch {owner}",
self.cfg.partition
);
}
Ok(self.stats)
}
async fn on_raw(&mut self, delivery: Delivery<RawBytes, RS::Handle>) {
self.stats.observed += 1;
let partition_class = Metrics::partition_class(self.cfg.partition);
self.metrics.observed(&partition_class);
if let Some(sent_at) = delivery.sent_at
&& let Ok(elapsed) = crate::bus::wallclock().duration_since(sent_at)
{
self.metrics
.raw_queue_wait_seconds(&partition_class, elapsed.as_secs_f64());
}
let now = Instant::now();
let reader = self.reader;
let envelope = RawEnvelope {
subject: &delivery.subject,
headers: Some(&delivery.headers),
bytes: delivery.item.0.as_slice(),
sent_at: delivery.sent_at,
handle: delivery.handle,
};
let disposition = reader.admit_bytes(&mut self.scheduler, &mut self.tracker, envelope, now);
match disposition {
RawDisposition::Queued { .. } => {
self.stats.queued += 1;
self.metrics.queued();
}
RawDisposition::Poison { handle, reason } => {
self.stats.poison += 1;
self.metrics.poison(reason.label());
debug!(reason = reason.label(), "poison raw message");
let seq = handle.sequence();
if handle.ack().await.is_ok() {
self.metrics.raw_acked("poison");
}
self.tracker.complete(seq);
}
RawDisposition::Suppressed { handle, reason } => {
self.stats.suppressed += 1;
self.metrics.suppressed(reason.label());
let seq = handle.sequence();
if handle.ack().await.is_ok() {
self.metrics.raw_acked("suppressed");
}
if !self.unpersisted.contains(&seq) {
self.tracker.complete(seq);
}
}
RawDisposition::Deferred { handle, reason } => {
self.stats.deferred += 1;
self.metrics.deferred(reason.label());
if reason == DeferReason::DuplicateOwned {
self.stats.coalesced += 1;
}
debug!(reason = reason.label(), "deferring raw message");
let _ = handle.nak(Some(self.cfg.blocked_retry)).await;
tokio::task::yield_now().await;
}
}
self.pump().await;
}
async fn on_result(&mut self, delivery: Delivery<SolveResult<E>, XS::Handle>) {
if let Some(sent_at) = delivery.sent_at
&& let Ok(elapsed) = crate::bus::wallclock().duration_since(sent_at)
{
self.metrics
.result_queue_wait_seconds(elapsed.as_secs_f64());
}
if let SolveOutcome::Solved { .. } = delivery.item.outcome
&& let Some(replica) = delivery.headers.get(STICKY_HEADER)
{
self.sticky.insert(
delivery.item.identity.vehicle_id,
(
replica.as_str().to_owned(),
Revision::from(delivery.item.identity.observation),
),
);
}
self.on_result_inner(delivery.item, delivery.handle).await;
self.pump().await;
}
async fn on_result_inner(&mut self, result: SolveResult<E>, handle: XS::Handle) {
let vehicle = result.identity.vehicle_id;
let now = Instant::now();
if let SolveOutcome::TripMiss = result.outcome {
self.resend_after_trip_miss(vehicle, result.job, now).await;
let _ = handle.ack().await;
return;
}
match self.result_verdict(vehicle, &result, now) {
ResultVerdict::Accept => {
self.requests.remove(&vehicle);
self.stats.accepted += 1;
let Some(meta) = self.active_meta.get(&vehicle).cloned() else {
let _ = handle.ack().await;
return;
};
self.metrics
.result(meta.region.as_str(), result.outcome.kind());
if let Some(active) = self.scheduler.active(vehicle) {
self.metrics.round_trip_seconds(
meta.region.as_str(),
result.outcome.kind(),
now.saturating_duration_since(active.dispatched)
.as_secs_f64(),
);
}
let result_for_retry = result.clone();
let decision = Decision::from_result(result, meta.reset, meta.segment);
self.commit_decision(
vehicle,
decision,
Some(ParkedResult {
result: result_for_retry,
handle,
}),
)
.await;
}
ResultVerdict::Park => {
let limit = self.cfg.parked_limit;
let tracked = self.scheduler.state(vehicle).is_some();
if !tracked || limit == 0 {
let _ = handle.nak(Some(self.cfg.blocked_retry)).await;
tokio::task::yield_now().await;
return;
}
let parked = self.parked_results.entry(vehicle).or_default();
if parked.len() < limit {
parked.push(ParkedResult { result, handle });
self.stats.parked += 1;
self.metrics.parked();
} else {
let _ = handle.nak(Some(self.cfg.blocked_retry)).await;
tokio::task::yield_now().await;
}
}
ResultVerdict::Reject(reason) => {
self.stats.rejected += 1;
self.metrics.rejected(reason.label());
debug!(
vehicle = vehicle.0,
reason = reason.label(),
"rejected result"
);
let _ = handle.ack().await;
}
ResultVerdict::Quarantine(reason) => {
self.stats.quarantined += 1;
self.metrics.quarantined(reason.label());
warn!(
vehicle = vehicle.0,
reason = reason.label(),
"quarantined solve result: kept out of the commit path"
);
if reason == QuarantineReason::DuplicateAfterCommit {
let _ = handle.nak(Some(self.cfg.blocked_retry)).await;
tokio::task::yield_now().await;
} else {
let _ = handle.ack().await;
}
}
}
}
async fn on_tick(&mut self) {
let now = Instant::now();
self.metrics.frontier_lag(
&Metrics::partition_class(self.cfg.partition),
self.tracker.outstanding() as u64,
);
self.metrics.oldest_pending_seconds(
self.scheduler
.oldest_pending(now)
.map_or(0.0, |age| age.as_secs_f64()),
);
let depth = self.scheduler.depth();
self.metrics.depth(
&Metrics::partition_class(self.cfg.partition),
depth.vehicles as u64,
depth.pending as u64,
depth.active as u64,
);
let heartbeat = self.epoch.is_some()
&& now.saturating_duration_since(self.last_heartbeat)
>= self.cfg.frontier.at_least_every;
let due = self.tracker.due(now).or_else(|| {
heartbeat.then_some(PartitionFrontier {
partition: self.cfg.partition,
sequence: self.tracker.frontier(),
})
});
if let Some(frontier) = due {
self.last_heartbeat = now;
match self.write_frontier(frontier).await {
Ok(true) => {
self.tracker.persisted(frontier.sequence, now);
self.stats.frontier = self.tracker.frontier();
}
Ok(false) => {
error!(
partition = self.cfg.partition,
epoch = ?self.epoch,
"frontier write fenced: a newer owner holds the partition"
);
self.fenced = Some(0);
}
Err(error) => {
warn!(partition = self.cfg.partition, %error, "frontier persist failed")
}
}
}
self.retry_blocked(now).await;
self.retry_requests(now).await;
self.persist_due(now);
if now.saturating_duration_since(self.last_evict) >= self.cfg.evict_every {
self.last_evict = now;
let dirty = &self.dirty;
let evicted = self
.scheduler
.evict_idle_except(now, |vehicle| dirty.contains_key(&vehicle));
for vehicle in evicted {
self.persisted_base.remove(&vehicle);
let _ = self.store.expire_idle(vehicle).await;
self.active_meta.remove(&vehicle);
self.requests.remove(&vehicle);
self.sticky.remove(&vehicle);
self.held.remove(&vehicle);
self.pending_reset.remove(&vehicle);
self.parked_results.remove(&vehicle);
}
}
self.retry_held().await;
self.pump().await;
}
async fn pump(&mut self) {
while let Some(vehicle) = self.scheduler.next_ready() {
if let DispatchOutcome::GloballyHeld = self.try_dispatch(vehicle).await {
break;
}
}
}
async fn try_dispatch(&mut self, vehicle: VehicleId) -> DispatchOutcome {
if self.quarantined.contains(&vehicle) {
return DispatchOutcome::Skipped;
}
if self.blocked.contains_key(&vehicle) {
return DispatchOutcome::Skipped;
}
if self.scheduler.active(vehicle).is_some() || self.scheduler.head(vehicle).is_none() {
return DispatchOutcome::Skipped;
}
let now = Instant::now();
if !self.scheduler.checkpoint(vehicle).is_loaded() {
let started = Instant::now();
let restored =
restore_vehicle(&self.store, vehicle, self.cfg.partition, &self.committer).await;
self.metrics
.checkpoint_restore_seconds(started.elapsed().as_secs_f64());
match restored {
Ok(restored) => {
if let Some(reason) = restored.reset {
self.pending_reset.insert(
vehicle,
PendingReset {
reason,
prior: restored.prior,
},
);
}
self.persisted_base.insert(vehicle, restored.prior);
self.scheduler.set_checkpoint(vehicle, restored.checkpoint);
for covered in self.scheduler.drain_committed(vehicle) {
let seq = covered.id.sequence;
self.stats.suppressed += 1;
self.metrics.suppressed("committed");
if covered.handle.ack().await.is_ok() {
self.metrics.raw_acked("suppressed");
}
self.tracker.complete(seq);
}
if self.scheduler.head(vehicle).is_none() {
return DispatchOutcome::Skipped;
}
}
Err(error) => {
warn!(vehicle = vehicle.0, %error, "restore failed; blocking vehicle");
self.blocked.insert(vehicle, now + self.cfg.blocked_retry);
return DispatchOutcome::Skipped;
}
}
}
let checkpoint = self.scheduler.checkpoint(vehicle).present().cloned();
let regression = self
.scheduler
.head(vehicle)
.and_then(|head| timestamp_regression(checkpoint.as_ref(), head.id, &head.payload));
if let Some(regression) = regression {
self.commit_timestamp_regression(vehicle, checkpoint.as_ref(), regression)
.await;
return DispatchOutcome::Done;
}
let head_point = self
.scheduler
.head(vehicle)
.expect("an eligible vehicle has a head")
.payload
.point;
let pin = checkpoint.as_ref().map(|cp| Pin {
region: cp.region.clone(),
graph: cp.graph.clone(),
routing_version: cp.routing_version,
});
let resolution = {
let resolver = Resolver::new(&self.catalog);
match resolver.resolve(vehicle, head_point, pin.as_ref()) {
Ok(resolution) => resolution,
Err(unserved) => {
self.commit_unserved(vehicle, checkpoint.as_ref(), unserved.cell)
.await;
return DispatchOutcome::Done;
}
}
};
let now_us = unix_micros();
let dispatch_started = Instant::now();
let sticky = checkpoint.as_ref().and_then(|cp| {
self.sticky
.get(&vehicle)
.filter(|(_, revision)| *revision == cp.revision)
.map(|(replica, _)| replica.as_str())
});
let dispatched = {
let head = self.scheduler.head(vehicle).expect("head");
self.dispatcher
.dispatch_to::<E, RS::Handle>(
vehicle,
head,
checkpoint.as_ref(),
&resolution,
&self.admission,
now,
now_us,
sticky,
)
.await
};
self.metrics.job_dispatch_seconds(
resolution.region.as_str(),
dispatch_started.elapsed().as_secs_f64(),
);
match dispatched {
Ok(d) => {
let (deadline, context) = if d.cached {
(self.cfg.sticky_retry, "cached")
} else {
(self.cfg.request_retry, "full")
};
self.metrics.solve_request(context, "dispatch");
self.requests.insert(vehicle, (d.request, now + deadline));
if self.metrics.sample_probe() {
self.metrics
.job_build_seconds(resolution.region.as_str(), d.built_for.as_secs_f64());
self.metrics.job_publish_seconds(
resolution.region.as_str(),
d.published_for.as_secs_f64(),
);
}
let pending = self.pending_reset.remove(&vehicle);
let reset = d.reset.or(pending.map(|p| p.reason));
let expected_base = d
.job
.identity
.base
.map(|b| b.revision)
.or_else(|| pending.and_then(|p| p.prior));
let job_bytes = d.job.bytes;
self.active_meta.insert(
vehicle,
ActiveMeta {
reset,
segment: d.segment,
region: resolution.region.clone(),
graph: resolution.graph.clone(),
routing_version: resolution.routing_version,
expected_base,
},
);
if self.scheduler.activate(vehicle, d.job).is_err() {
self.active_meta.remove(&vehicle);
self.requests.remove(&vehicle);
return DispatchOutcome::Skipped;
}
self.held.remove(&vehicle);
self.stats.dispatched += 1;
self.metrics
.dispatched(resolution.region.as_str(), resolution.lane.0);
self.metrics
.job_bytes(resolution.region.as_str(), job_bytes);
if let Some(parked) = self.parked_results.remove(&vehicle) {
for parked in parked {
self.on_result_inner(parked.result, parked.handle).await;
}
}
DispatchOutcome::Done
}
Err(DispatchError::Held(reason)) => {
self.stats.held += 1;
self.metrics.held(resolution.region.as_str());
debug!(vehicle = vehicle.0, reason = %reason, "dispatch held by admission");
self.held
.insert(vehicle, self.admission.hold(&resolution.region));
match reason {
HeldReason::RegionJobs | HeldReason::RegionBytes => DispatchOutcome::RegionHeld,
HeldReason::GlobalJobs | HeldReason::GlobalBytes => {
DispatchOutcome::GloballyHeld
}
}
}
Err(DispatchError::TimestampRegression(regression)) => {
self.commit_timestamp_regression(vehicle, checkpoint.as_ref(), regression)
.await;
DispatchOutcome::Done
}
Err(error) => {
warn!(vehicle = vehicle.0, %error, "dispatch failed; will retry");
self.held
.insert(vehicle, self.admission.hold(&resolution.region));
DispatchOutcome::GloballyHeld
}
}
}
async fn commit_timestamp_regression(
&mut self,
vehicle: VehicleId,
checkpoint: Option<&VehicleCheckpoint<E>>,
regression: TimestampRegression,
) {
let Some(checkpoint) = checkpoint else {
debug_assert!(false, "a timestamp regression requires a checkpoint");
return;
};
debug!(
vehicle = vehicle.0,
committed_timestamp = regression.committed,
incoming_timestamp = regression.incoming,
"committing timestamp-regression terminal"
);
self.commit_local_terminal(
vehicle,
LocalTerminalRoute {
region: checkpoint.region.clone(),
graph: checkpoint.graph.clone(),
routing_version: checkpoint.routing_version,
},
Some(checkpoint),
TerminalReason::TimestampRegression,
JobReservation::SyntheticTerminal,
)
.await;
}
async fn commit_decision(
&mut self,
vehicle: VehicleId,
decision: Decision<E>,
result_delivery: Option<ParkedResult<E, XS::Handle>>,
) {
let Some(meta) = self.active_meta.get(&vehicle).cloned() else {
if let Some(delivery) = result_delivery {
let _ = delivery.handle.ack().await;
}
return;
};
if self.scheduler.begin_commit(vehicle).is_err() {
if let Some(delivery) = result_delivery {
let _ = delivery.handle.ack().await;
}
return;
}
let raw = self
.scheduler
.active(vehicle)
.expect("begin_commit succeeded, so a job is active")
.observation;
let prev = self.scheduler.checkpoint(vehicle).present().cloned();
let (reset_opt, segment) = decision_reset_segment(&decision);
debug_assert!(
reset_opt.is_some() || prev.as_ref().is_none_or(|p| p.segment == segment),
"a commit must not change the segment without a reset",
);
let is_terminal = matches!(decision, Decision::Terminal { .. });
let terminal_reason = match &decision {
Decision::Terminal { reason, .. } => Some(reason.label()),
_ => None,
};
let reset_reason = reset_opt.map(|reason| reason.label());
let plan = commit::plan(
prev.as_ref(),
decision,
raw,
&meta.region,
&meta.graph,
meta.routing_version,
);
let next = plan.next.clone();
let committer = self.committer.clone();
let partition = self.cfg.partition;
let expected_base = meta.expected_base;
let sampled = self.metrics.sample_probe();
let started = Instant::now();
let guard = self.drain.begin();
self.commits.push(Box::pin(async move {
let result = committer
.commit_sampled(vehicle, partition, plan, expected_base, raw, sampled)
.await;
drop(guard);
CommitCompletion {
vehicle,
next,
result_delivery,
started,
completed: Instant::now(),
sampled,
is_terminal,
terminal_reason,
reset_reason,
result,
}
}));
}
async fn finish_commit(&mut self, completion: CommitCompletion<E, XS::Handle, S::Error>) {
let CommitCompletion {
vehicle,
next,
result_delivery,
started,
completed,
sampled,
is_terminal,
terminal_reason,
reset_reason,
result,
} = completion;
let now = Instant::now();
if sampled {
self.metrics
.commit_stage_seconds("handoff", now.duration_since(completed).as_secs_f64());
}
match result {
Ok(committed) => {
self.scheduler
.set_checkpoint(vehicle, CheckpointState::Present(next));
if let Ok(finished) = self.scheduler.finish(vehicle, now) {
let seq = finished.observation.id.sequence;
let raw_ack_started = Instant::now();
let raw_acked = finished.observation.handle.ack().await;
if sampled {
self.metrics.commit_stage_seconds(
"raw_ack",
raw_ack_started.elapsed().as_secs_f64(),
);
}
if raw_acked.is_ok() {
self.metrics.raw_acked("committed");
}
if let Some(delivery) = result_delivery {
let result_ack_started = Instant::now();
let _ = delivery.handle.ack().await;
if sampled {
self.metrics.commit_stage_seconds(
"result_ack",
result_ack_started.elapsed().as_secs_f64(),
);
}
}
if self.cfg.commit.mode == CommitMode::Deferred {
self.defer_completion(vehicle, seq, now);
} else {
self.tracker.complete(seq);
}
drop(finished.job); }
self.active_meta.remove(&vehicle);
self.requests.remove(&vehicle);
self.pending_reset.remove(&vehicle);
self.blocked.remove(&vehicle);
self.stats.committed += 1;
let kind = if is_terminal { "terminal" } else { "matched" };
self.metrics
.commit_seconds(kind, started.elapsed().as_secs_f64());
self.metrics.output_bytes(committed.bytes as u64);
if let Some(reason) = reset_reason {
self.metrics.completion("reset", reason);
}
if is_terminal {
self.metrics
.completion("terminal", terminal_reason.unwrap_or("internal"));
} else {
self.metrics.completion("matched", "solved");
}
if is_terminal {
self.stats.terminal += 1;
}
if reset_reason.is_some() {
self.stats.resets += 1;
}
self.stats.frontier = self.tracker.frontier();
}
Err(CommitError::Conflict { actual }) => {
self.stats.conflicts += 1;
warn!(
vehicle = vehicle.0,
?actual,
"commit conflicted; blocking vehicle"
);
if let Ok(restored) =
restore_vehicle(&self.store, vehicle, self.cfg.partition, &self.committer).await
{
self.persisted_base.insert(vehicle, restored.prior);
self.scheduler.set_checkpoint(vehicle, restored.checkpoint);
}
if let Some(delivery) = result_delivery {
self.blocked_results.insert(vehicle, delivery);
}
self.blocked.insert(vehicle, now + self.cfg.blocked_retry);
}
Err(CommitError::Fenced { owner }) => {
error!(vehicle = vehicle.0, owner, "commit fenced by a newer owner");
self.fenced = Some(owner);
}
Err(error) if is_permanent_commit_error(&error) => {
self.quarantine_vehicle(vehicle, &error);
}
Err(error) => {
warn!(vehicle = vehicle.0, %error, "commit incomplete; blocking vehicle");
if let Some(delivery) = result_delivery {
self.blocked_results.insert(vehicle, delivery);
}
self.blocked.insert(vehicle, now + self.cfg.blocked_retry);
}
}
}
async fn write_frontier(&mut self, frontier: PartitionFrontier) -> Result<bool, S::Error> {
match self.epoch {
Some(epoch) => self.store.set_frontier_fenced(frontier, epoch).await,
None => self.store.set_frontier(frontier).await.map(|()| true),
}
}
fn defer_completion(&mut self, vehicle: VehicleId, seq: u64, now: Instant) {
let dirty = self.dirty.entry(vehicle).or_insert_with(|| Dirty {
sequences: Vec::new(),
since: now,
in_flight: false,
retry_at: None,
});
if dirty.sequences.is_empty() {
dirty.since = now;
}
dirty.sequences.push(seq);
self.unpersisted.insert(seq);
if dirty.sequences.len() >= self.cfg.checkpoint.every {
self.start_persist(vehicle, now, false);
}
}
fn persist_due(&mut self, now: Instant) {
if self.dirty.is_empty() {
return;
}
let policy = self.cfg.checkpoint;
let due: Vec<VehicleId> = self
.dirty
.iter()
.filter(|(_, dirty)| {
!dirty.in_flight
&& dirty.retry_at.is_none_or(|at| at <= now)
&& (dirty.sequences.len() >= policy.every
|| now.saturating_duration_since(dirty.since) >= policy.interval)
})
.map(|(&vehicle, _)| vehicle)
.collect();
for vehicle in due {
if !self.start_persist(vehicle, now, false) {
break;
}
}
}
fn start_persist(&mut self, vehicle: VehicleId, now: Instant, unbounded: bool) -> bool {
if !unbounded && self.persists.len() >= self.cfg.checkpoint.in_flight {
return false;
}
let Some(dirty) = self.dirty.get_mut(&vehicle) else {
return true;
};
if dirty.in_flight || dirty.sequences.is_empty() {
return true;
}
let Some(checkpoint) = self.scheduler.checkpoint(vehicle).present().cloned() else {
return true;
};
dirty.in_flight = true;
dirty.retry_at = None;
let sequences = core::mem::take(&mut dirty.sequences);
let base = self.persisted_base.get(&vehicle).copied().flatten();
let committer = self.committer.clone();
let revision = checkpoint.revision;
let guard = self.drain.begin();
self.persists.push(Box::pin(async move {
let result = committer.persist(vehicle, base, &checkpoint).await;
drop(guard);
PersistCompletion {
vehicle,
revision,
sequences,
started: now,
result,
}
}));
true
}
async fn finish_persist(&mut self, completion: PersistCompletion<S::Error>) {
let PersistCompletion {
vehicle,
revision,
sequences,
started,
result,
} = completion;
let now = Instant::now();
self.metrics.commit_stage_seconds(
"checkpoint_persist",
now.saturating_duration_since(started).as_secs_f64(),
);
let mut adopted = false;
let persisted = match result {
Ok(DirectOutcome::Committed | DirectOutcome::AlreadyCommitted) => true,
Ok(DirectOutcome::Conflict { actual }) => {
adopted = self.adopt_stored_base(vehicle, revision, actual, now).await;
false
}
Ok(DirectOutcome::Busy { pending }) => {
warn!(vehicle = vehicle.0, %pending, "deferred checkpoint waits on a staged record");
false
}
Ok(DirectOutcome::Fenced { owner }) => {
error!(
vehicle = vehicle.0,
owner, "deferred checkpoint fenced by a newer owner"
);
self.fenced = Some(owner);
false
}
Err(error) => {
warn!(vehicle = vehicle.0, %error, "deferred checkpoint persist failed; will retry");
false
}
};
if persisted {
self.stats.persisted += 1;
self.persisted_base.insert(vehicle, Some(revision));
for seq in &sequences {
self.unpersisted.remove(seq);
self.tracker.complete(*seq);
}
self.stats.frontier = self.tracker.frontier();
}
let Some(dirty) = self.dirty.get_mut(&vehicle) else {
return;
};
dirty.in_flight = false;
if !persisted {
let newer = core::mem::replace(&mut dirty.sequences, sequences);
dirty.sequences.extend(newer);
dirty.retry_at = Some(if adopted {
now
} else {
now + self.cfg.blocked_retry
});
} else if dirty.sequences.is_empty() {
self.dirty.remove(&vehicle);
} else if dirty.sequences.len() >= self.cfg.checkpoint.every {
self.start_persist(vehicle, now, false);
}
}
async fn adopt_stored_base(
&mut self,
vehicle: VehicleId,
revision: Revision,
actual: Option<Revision>,
now: Instant,
) -> bool {
let behind = actual.is_none_or(|stored| stored <= revision);
if self.epoch.is_none() || !behind {
warn!(
vehicle = vehicle.0,
?actual,
?revision,
"deferred checkpoint lost its base; frontier held for replay"
);
return false;
}
let frontier = PartitionFrontier {
partition: self.cfg.partition,
sequence: self.tracker.frontier(),
};
match self.write_frontier(frontier).await {
Ok(true) => {
self.last_heartbeat = now;
self.tracker.persisted(frontier.sequence, now);
warn!(
vehicle = vehicle.0,
?actual,
?revision,
"an earlier owner saved this vehicle; adopting its revision as the base"
);
self.persisted_base.insert(vehicle, actual);
true
}
Ok(false) => {
error!(
partition = self.cfg.partition,
epoch = ?self.epoch,
"frontier write fenced: a newer owner holds the partition"
);
self.fenced = Some(0);
false
}
Err(error) => {
warn!(vehicle = vehicle.0, %error, "ownership check failed; will retry the persist");
false
}
}
}
fn quarantine_vehicle(&mut self, vehicle: VehicleId, error: &CommitError<S::Error>) {
error!(
vehicle = vehicle.0,
%error,
"commit failed permanently; quarantining vehicle for this process lifetime"
);
self.blocked.remove(&vehicle);
self.pending_reset.remove(&vehicle);
self.held.remove(&vehicle);
if self.quarantined.insert(vehicle) {
self.stats.quarantined_vehicles += 1;
}
}
async fn commit_unserved(
&mut self,
vehicle: VehicleId,
checkpoint: Option<&VehicleCheckpoint<E>>,
cell: String,
) {
let (region, graph, routing_version) = match checkpoint {
Some(cp) => (cp.region.clone(), cp.graph.clone(), cp.routing_version),
None => match self.catalog.regions.first() {
Some(region) => (
region.id.clone(),
region.graph.clone(),
self.catalog.routing_version,
),
None => {
warn!(vehicle = vehicle.0, %cell, "unserved cell and empty catalog; cannot advance");
return;
}
},
};
debug!(vehicle = vehicle.0, %cell, "committing unsupported-coverage terminal");
self.commit_local_terminal(
vehicle,
LocalTerminalRoute {
region,
graph,
routing_version,
},
checkpoint,
TerminalReason::UnsupportedCoverage,
JobReservation::SyntheticTerminal,
)
.await;
}
async fn commit_local_terminal(
&mut self,
vehicle: VehicleId,
route: LocalTerminalRoute,
checkpoint: Option<&VehicleCheckpoint<E>>,
reason: TerminalReason,
reservation: JobReservation,
) {
let LocalTerminalRoute {
region,
graph,
routing_version,
} = route;
let now = Instant::now();
let head_obs = self
.scheduler
.head(vehicle)
.expect("an eligible vehicle has a head")
.id;
let segment = checkpoint.map_or(SegmentId::from(head_obs), |cp| cp.segment);
let base = checkpoint.map(|cp| BaseState {
revision: cp.revision,
segment: cp.segment,
});
let identity = JobIdentity {
schema: SCHEMA_VERSION,
vehicle_id: vehicle,
observation: head_obs,
base,
graph: graph.clone(),
region: region.clone(),
};
let job_id = identity.local_decision_id();
let expected_base = base
.map(|b| b.revision)
.or_else(|| self.pending_reset.get(&vehicle).and_then(|p| p.prior));
self.active_meta.insert(
vehicle,
ActiveMeta {
reset: None,
segment,
region,
graph,
routing_version,
expected_base,
},
);
let job = ActiveJob {
id: job_id,
identity: identity.clone(),
observation: head_obs,
bytes: 0,
reservation,
dispatched: now,
};
if self.scheduler.activate(vehicle, job).is_err() {
self.active_meta.remove(&vehicle);
return;
}
self.pending_reset.remove(&vehicle);
let decision = Decision::Terminal {
job: job_id,
identity,
reason,
closes_segment: false,
segment,
};
self.commit_decision(vehicle, decision, None).await;
}
async fn retry_blocked(&mut self, now: Instant) {
let due: Vec<VehicleId> = self
.blocked
.iter()
.filter(|&(_, &at)| at <= now)
.map(|(&vehicle, _)| vehicle)
.collect();
for vehicle in due {
if self.cfg.commit.mode == CommitMode::Deferred {
self.retry_unprepared(vehicle, now).await;
continue;
}
let prepared = match self.store.load(vehicle).await {
Ok((_, prepared)) => prepared,
Err(error) => {
warn!(vehicle = vehicle.0, %error, "blocked-vehicle load failed");
self.blocked.insert(vehicle, now + self.cfg.blocked_retry);
continue;
}
};
match prepared {
Some(prepared) => {
match self
.committer
.finish_prepared(vehicle, self.cfg.partition, prepared)
.await
{
Ok(committed) => {
let active = self.scheduler.active(vehicle).map(|job| job.observation);
if finished_prepared_is_active(active, committed.raw) {
self.resolve_blocked(vehicle, now).await;
} else {
self.retry_unprepared(vehicle, now).await;
}
}
Err(error) if is_permanent_commit_error(&error) => {
self.quarantine_vehicle(vehicle, &error);
}
Err(error) => {
debug!(vehicle = vehicle.0, %error, "blocked commit still cannot finish");
self.blocked.insert(vehicle, now + self.cfg.blocked_retry);
}
}
}
None => self.retry_unprepared(vehicle, now).await,
}
}
}
async fn retry_unprepared(&mut self, vehicle: VehicleId, now: Instant) {
let Some(observation) = self.scheduler.active(vehicle).map(|job| job.observation) else {
self.active_meta.remove(&vehicle);
self.blocked.remove(&vehicle);
return;
};
if self.cfg.commit.mode == CommitMode::Deferred {
self.retry_from_memory(vehicle);
return;
}
let restored = match restore_vehicle(
&self.store,
vehicle,
self.cfg.partition,
&self.committer,
)
.await
{
Ok(restored) => restored,
Err(error) => {
debug!(vehicle = vehicle.0, %error, "blocked commit still cannot restore");
self.blocked.insert(vehicle, now + self.cfg.blocked_retry);
return;
}
};
let committed = restored
.checkpoint
.present()
.is_some_and(|checkpoint| checkpoint.last_input >= observation);
self.scheduler.set_checkpoint(vehicle, restored.checkpoint);
if committed {
self.resolve_blocked(vehicle, now).await;
return;
}
if let Some(reason) = restored.reset {
self.pending_reset.insert(
vehicle,
PendingReset {
reason,
prior: restored.prior,
},
);
}
self.retry_from_memory(vehicle);
}
fn retry_from_memory(&mut self, vehicle: VehicleId) {
if let Some(delivery) = self.blocked_results.remove(&vehicle) {
self.parked_results
.entry(vehicle)
.or_default()
.push(delivery);
}
if let Some(job) = self.scheduler.rollback_commit(vehicle) {
drop(job);
}
self.active_meta.remove(&vehicle);
self.requests.remove(&vehicle);
self.blocked.remove(&vehicle);
self.held.remove(&vehicle);
}
async fn resolve_blocked(&mut self, vehicle: VehicleId, now: Instant) {
if self.scheduler.active(vehicle).is_some() {
self.scheduler
.set_checkpoint(vehicle, CheckpointState::Unloaded);
if let Ok(finished) = self.scheduler.finish(vehicle, now) {
let seq = finished.observation.id.sequence;
if finished.observation.handle.ack().await.is_ok() {
self.metrics.raw_acked("committed");
}
self.tracker.complete(seq);
drop(finished.job);
}
}
self.active_meta.remove(&vehicle);
self.blocked.remove(&vehicle);
if let Some(delivery) = self.blocked_results.remove(&vehicle) {
let _ = delivery.handle.ack().await;
}
self.stats.frontier = self.tracker.frontier();
}
async fn retry_requests(&mut self, now: Instant) {
let due: Vec<VehicleId> = self
.requests
.iter()
.filter(|(_, (_, at))| *at <= now)
.map(|(&vehicle, _)| vehicle)
.collect();
for vehicle in due {
if self.blocked.contains_key(&vehicle) || self.scheduler.active(vehicle).is_none() {
continue;
}
self.sticky.remove(&vehicle);
self.metrics.solve_request("full", "timeout");
let Some((request, at)) = self.requests.get_mut(&vehicle) else {
continue;
};
if let Err(error) = self.dispatcher.republish::<E>(request).await {
warn!(vehicle = vehicle.0, %error, "transient solve request retry failed");
}
*at = Instant::now() + self.cfg.request_retry;
}
}
async fn resend_after_trip_miss(&mut self, vehicle: VehicleId, job: JobId, now: Instant) {
self.sticky.remove(&vehicle);
let active = self
.scheduler
.active(vehicle)
.is_some_and(|active| active.id == job);
if !active || self.blocked.contains_key(&vehicle) {
return;
}
self.metrics.solve_request("full", "trip_miss");
let Some((request, at)) = self.requests.get_mut(&vehicle) else {
return;
};
if let Err(error) = self.dispatcher.republish::<E>(request).await {
warn!(vehicle = vehicle.0, %error, "full resend after trip miss failed");
}
*at = now + self.cfg.request_retry;
}
async fn retry_held(&mut self) {
let held: Vec<VehicleId> = self.held.keys().copied().collect();
for vehicle in held {
if self.blocked.contains_key(&vehicle) {
continue;
}
if let DispatchOutcome::GloballyHeld = self.try_dispatch(vehicle).await {
break;
}
}
}
fn result_verdict(
&self,
vehicle: VehicleId,
result: &SolveResult<E>,
now: Instant,
) -> ResultVerdict {
match self.scheduler.state(vehicle) {
Some(state) => validate::validate(state, result, self.cfg.partition).into(),
None => {
let detached = detached_state::<E, XS::Handle>(now);
validate::validate(&detached, result, self.cfg.partition).into()
}
}
}
}
fn is_permanent_commit_error<SE>(error: &CommitError<SE>) -> bool {
matches!(error, CommitError::Decode(_) | CommitError::Encode(_))
}
fn finished_prepared_is_active(
active: Option<crate::protocol::ids::ObservationId>,
committed_raw: crate::protocol::ids::ObservationId,
) -> bool {
active == Some(committed_raw)
}
fn detached_state<E: Entry, H: AckHandle>(now: Instant) -> VehicleState<E, H> {
VehicleState {
checkpoint: CheckpointState::Unloaded,
pending: VecDeque::new(),
active: None,
committing: false,
last_touch: now,
}
}
fn decision_reset_segment<E: Entry>(decision: &Decision<E>) -> (Option<ResetReason>, SegmentId) {
match decision {
Decision::Solved { reset, segment, .. } => (*reset, *segment),
Decision::Terminal { segment, .. } => (None, *segment),
Decision::Reset {
reason,
new_segment,
..
} => (Some(*reason), *new_segment),
}
}
fn unix_micros() -> i64 {
crate::bus::wallclock()
.duration_since(UNIX_EPOCH)
.ok()
.and_then(|elapsed| i64::try_from(elapsed.as_micros()).ok())
.unwrap_or(i64::MAX)
}
#[cfg(test)]
mod tests {
use super::*;
use core::sync::atomic::{AtomicUsize, Ordering};
use async_nats::HeaderMap;
use chrono::{DateTime, Utc};
use geo::Point;
use routers_network::mock::MockEntryId;
use routers_transition::matcher::{Origin, Trip};
use tokio::sync::Notify;
use crate::bus::Wire;
use crate::bus::memory::{MemoryBus, MemoryPublisher, MemorySource};
use crate::event::{MatchedDiff, Payload, shard_of};
use crate::orchestrator::admission::AdmissionConfig;
use crate::partition::partition_of;
use crate::protocol::ids::{ObservationId, OutputId};
use crate::protocol::job::SolveJob;
use crate::protocol::job::SolveRequest;
use crate::protocol::output::OutputKind;
use crate::protocol::result::SolveOutcome;
use crate::store::checkpoint::{CommitPhase, MemoryCheckpointStore, PreparedCommit};
use crate::topology::{output_subject, raw_subject, reply_subject};
const REPLY: &str = "_INBOX.test";
const REQUESTS: &str = "solve.req.v1.>";
fn full_job(bytes: &[u8]) -> SolveJob<E> {
SolveRequest::<E>::decode(bytes)
.expect("request decodes")
.resolve(|_, _| None)
.expect("a full request resolves")
}
fn results_subject(partition: u16) -> String {
reply_subject(REPLY, partition)
}
type E = MockEntryId;
type Worker = PartitionWorker<
E,
MemoryCheckpointStore,
MemoryPublisher<SolveRequest<E>>,
MemoryPublisher<CommittedOutput<E>>,
MemorySource<RawBytes>,
MemorySource<SolveResult<E>>,
>;
struct OverlapPublisher<T: Wire> {
inner: MemoryPublisher<T>,
active: Arc<AtomicUsize>,
peak: Arc<AtomicUsize>,
changed: Arc<Notify>,
}
impl<T: Wire> Clone for OverlapPublisher<T> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
active: Arc::clone(&self.active),
peak: Arc::clone(&self.peak),
changed: Arc::clone(&self.changed),
}
}
}
impl<T: Wire> OverlapPublisher<T> {
fn new(inner: MemoryPublisher<T>) -> Self {
Self {
inner,
active: Arc::new(AtomicUsize::new(0)),
peak: Arc::new(AtomicUsize::new(0)),
changed: Arc::new(Notify::new()),
}
}
}
impl<T: Wire + Send + Sync + 'static> Publisher<T> for OverlapPublisher<T> {
async fn publish_bytes(
&self,
subject: &str,
msg_id: &str,
headers: HeaderMap,
bytes: &[u8],
) -> Result<crate::bus::adapter::PublishOutcome, crate::bus::adapter::PublishError>
{
let active = self.active.fetch_add(1, Ordering::SeqCst) + 1;
self.peak.fetch_max(active, Ordering::SeqCst);
if active == 1 {
while self.peak.load(Ordering::SeqCst) < 2 {
self.changed.notified().await;
}
} else {
self.changed.notify_waiters();
}
let result = self
.inner
.publish_bytes(subject, msg_id, headers, bytes)
.await;
self.active.fetch_sub(1, Ordering::SeqCst);
result
}
}
fn point() -> Point {
Point::new(151.2093, -33.8688)
}
fn catalog(budget_ms: u64) -> Catalog {
let cell = shard_of(point()).to_string();
let toml = format!(
r#"
version = 1
routing_version = 1
[[regions]]
id = "r1"
graph = "g1"
coverage = ["{cell}"]
overlap = []
lanes = 1
resource_class = "c"
replicas = {{ min = 1, max = 1 }}
freshness_budget_ms = {budget_ms}
"#,
);
Catalog::parse(&toml).expect("catalog parses")
}
fn second_point() -> Point {
Point::new(144.9631, -37.8136)
}
fn two_region_catalog() -> Catalog {
let first = shard_of(point()).to_string();
let second = shard_of(second_point()).to_string();
assert_ne!(first, second);
let toml = format!(
r#"
version = 1
routing_version = 1
[[regions]]
id = "r1"
graph = "g1"
coverage = ["{first}"]
overlap = []
lanes = 1
resource_class = "c"
replicas = {{ min = 1, max = 1 }}
freshness_budget_ms = 30000
[[regions]]
id = "r2"
graph = "g2"
coverage = ["{second}"]
overlap = []
lanes = 1
resource_class = "c"
replicas = {{ min = 1, max = 1 }}
freshness_budget_ms = 30000
"#,
);
Catalog::parse(&toml).expect("catalog parses")
}
fn config(partition: u16) -> WorkerConfig {
WorkerConfig {
tick: Duration::from_millis(20),
blocked_retry: Duration::from_millis(20),
..WorkerConfig::new(partition)
}
}
fn admission_for(catalog: &Catalog) -> Admission {
Admission::new(
AdmissionConfig::default(),
catalog.regions.iter().map(|r| &r.id),
)
}
fn clean_report(partition: u16) -> RecoveryReport {
RecoveryReport {
partition,
frontier: None,
prepared_found: 0,
prepared_finished: 0,
prepared_failed: Vec::new(),
epoch: None,
}
}
fn build_worker(
cfg: WorkerConfig,
catalog: Arc<Catalog>,
bus: &MemoryBus,
store: &MemoryCheckpointStore,
shutdown: Shutdown,
) -> Worker {
let admission = admission_for(&catalog);
let report = clean_report(cfg.partition);
build_worker_with(cfg, catalog, bus, store, admission, &report, shutdown)
}
#[allow(clippy::too_many_arguments)]
fn build_worker_with(
cfg: WorkerConfig,
catalog: Arc<Catalog>,
bus: &MemoryBus,
store: &MemoryCheckpointStore,
admission: Admission,
report: &RecoveryReport,
shutdown: Shutdown,
) -> Worker {
let dispatcher = Dispatcher::new(
bus.publisher::<SolveRequest<E>>(),
DispatchConfig::default(),
REPLY.to_owned(),
);
let committer = Committer::new(
store.clone(),
bus.publisher::<CommittedOutput<E>>(),
CommitConfig {
publish_attempts: 4,
backoff: Duration::ZERO,
..CommitConfig::default()
},
);
let raw = bus.source::<RawBytes>(&raw_subject(u64::from(cfg.partition)));
let results = bus.source::<SolveResult<E>>(&results_subject(cfg.partition));
PartitionWorker::new(
cfg,
catalog,
admission,
store.clone(),
dispatcher,
committer,
raw,
results,
report,
shutdown,
)
}
fn partition_for(vehicle: u64) -> u16 {
partition_of(VehicleId(vehicle)) as u16
}
#[test]
fn foreign_prepared_completion_cannot_release_the_active_raw() {
let active = ObservationId {
partition: 7,
sequence: 12,
};
let foreign = ObservationId {
partition: 7,
sequence: 11,
};
assert!(finished_prepared_is_active(Some(active), active));
assert!(!finished_prepared_is_active(Some(active), foreign));
assert!(!finished_prepared_is_active(None, active));
}
#[tokio::test]
async fn raw_owned_request_is_retained_until_answered_and_retried_on_timeout() {
let vehicle = 1_u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let catalog = Arc::new(two_region_catalog());
let mut cfg = config(partition);
cfg.request_retry = Duration::from_millis(1);
let mut worker = build_worker(cfg, catalog, &bus, &store, Shutdown::new());
publish_raw_at(&bus, partition, vehicle, 1_775_000_000_000_000, point()).await;
enqueue_next_raw(&mut worker).await;
worker.pump().await;
assert!(worker.requests.contains_key(&VehicleId(vehicle)));
let (request, deadline) = worker.requests.get_mut(&VehicleId(vehicle)).unwrap();
assert_eq!(request.reply, results_subject(partition));
*deadline = Instant::now();
worker.retry_requests(Instant::now()).await;
assert!(worker.requests[&VehicleId(vehicle)].1 > Instant::now());
assert_eq!(
bus.published(REQUESTS).len(),
1,
"memory bus deduplicates the exact retry"
);
}
#[tokio::test]
async fn pump_continues_after_a_region_local_admission_hold() {
let first = 1u64;
let second = same_partition_as(first);
let partition = partition_for(first);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
let catalog = Arc::new(two_region_catalog());
let admission = Admission::new(
AdmissionConfig {
global_jobs: 10,
region_jobs: 1,
..AdmissionConfig::default()
},
catalog.regions.iter().map(|region| ®ion.id),
);
let occupied = admission
.try_admit(&RegionId::new("r1").unwrap(), 1)
.expect("r1 has one slot");
publish_raw_at(&bus, partition, first, 1_775_000_000_000_000, point()).await;
publish_raw_at(
&bus,
partition,
second,
1_775_000_000_000_000,
second_point(),
)
.await;
let report = clean_report(partition);
let mut worker = build_worker_with(
config(partition),
catalog,
&bus,
&store,
admission,
&report,
shutdown,
);
enqueue_next_raw(&mut worker).await;
enqueue_next_raw(&mut worker).await;
worker.pump().await;
let published = bus.published(REQUESTS);
assert_eq!(
published.len(),
1,
"the free region dispatched in this pass"
);
let job = full_job(&published[0].2);
assert_eq!(job.identity.vehicle_id, VehicleId(second));
assert_eq!(job.identity.region, RegionId::new("r2").unwrap());
assert!(worker.held.contains_key(&VehicleId(first)));
drop(occupied);
}
#[tokio::test]
async fn pump_stops_after_a_global_admission_hold() {
let first = 1u64;
let second = same_partition_as(first);
let partition = partition_for(first);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
let catalog = Arc::new(two_region_catalog());
let admission = Admission::new(
AdmissionConfig {
global_jobs: 1,
region_jobs: 10,
..AdmissionConfig::default()
},
catalog.regions.iter().map(|region| ®ion.id),
);
let occupied = admission
.try_admit(&RegionId::new("r1").unwrap(), 1)
.expect("the process has one slot");
publish_raw_at(&bus, partition, first, 1_775_000_000_000_000, point()).await;
publish_raw_at(
&bus,
partition,
second,
1_775_000_000_000_000,
second_point(),
)
.await;
let report = clean_report(partition);
let mut worker = build_worker_with(
config(partition),
catalog,
&bus,
&store,
admission,
&report,
shutdown,
);
enqueue_next_raw(&mut worker).await;
enqueue_next_raw(&mut worker).await;
worker.pump().await;
assert!(
bus.published(REQUESTS).is_empty(),
"global pressure stops the round before another region is attempted"
);
assert!(worker.held.contains_key(&VehicleId(first)));
assert!(!worker.held.contains_key(&VehicleId(second)));
drop(occupied);
}
#[tokio::test]
async fn retry_held_continues_past_a_still_saturated_region() {
let first = 1u64;
let second = same_partition_as(first);
let partition = partition_for(first);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
let catalog = Arc::new(two_region_catalog());
let admission = Admission::new(
AdmissionConfig {
global_jobs: 10,
region_jobs: 1,
..AdmissionConfig::default()
},
catalog.regions.iter().map(|region| ®ion.id),
);
publish_raw_at(&bus, partition, first, 1_775_000_000_000_000, point()).await;
publish_raw_at(
&bus,
partition,
second,
1_775_000_000_000_000,
second_point(),
)
.await;
let report = clean_report(partition);
let mut worker = build_worker_with(
config(partition),
catalog,
&bus,
&store,
admission.clone(),
&report,
shutdown,
);
enqueue_next_raw(&mut worker).await;
enqueue_next_raw(&mut worker).await;
worker.held.insert(
VehicleId(first),
admission.hold(&RegionId::new("r1").unwrap()),
);
worker.held.insert(
VehicleId(second),
admission.hold(&RegionId::new("r2").unwrap()),
);
let first_retry = *worker.held.keys().next().expect("two held vehicles");
let (blocked_region, expected_vehicle) = if first_retry == VehicleId(first) {
(RegionId::new("r1").unwrap(), VehicleId(second))
} else {
(RegionId::new("r2").unwrap(), VehicleId(first))
};
let occupied = admission
.try_admit(&blocked_region, 1)
.expect("the first region has one slot");
worker.retry_held().await;
let published = bus.published(REQUESTS);
assert_eq!(published.len(), 1, "the free region retried in this pass");
let job = full_job(&published[0].2);
assert_eq!(job.identity.vehicle_id, expected_vehicle);
assert!(worker.held.contains_key(&first_retry));
assert!(!worker.held.contains_key(&expected_vehicle));
drop(occupied);
}
fn same_partition_as(vehicle: u64) -> u64 {
let want = partition_for(vehicle);
(vehicle + 1..)
.find(|&v| partition_for(v) == want)
.expect("another vehicle in the partition")
}
async fn publish_raw_at(
bus: &MemoryBus,
partition: u16,
vehicle: u64,
ts_us: i64,
point: Point,
) -> u64 {
let payload = Payload {
vehicle_id: VehicleId(vehicle),
timestamp: DateTime::<Utc>::from_timestamp_micros(ts_us).unwrap(),
point,
};
let bytes = payload.encode().unwrap();
let msg_id = format!("{vehicle}:{ts_us}");
let outcome = bus
.publisher::<RawBytes>()
.publish_bytes(
&raw_subject(u64::from(partition)),
&msg_id,
crate::bus::outbound(),
&bytes,
)
.await
.expect("raw publish");
match outcome {
crate::bus::adapter::PublishOutcome::Acked { sequence, .. } => sequence,
}
}
async fn publish_raw(bus: &MemoryBus, partition: u16, vehicle: u64, ts_us: i64) -> u64 {
publish_raw_at(bus, partition, vehicle, ts_us, point()).await
}
async fn enqueue_next_raw<OP>(
worker: &mut PartitionWorker<
E,
MemoryCheckpointStore,
MemoryPublisher<SolveRequest<E>>,
OP,
MemorySource<RawBytes>,
MemorySource<SolveResult<E>>,
>,
) where
OP: Publisher<CommittedOutput<E>>,
{
let delivery = worker
.raw
.next()
.await
.expect("raw source stays open")
.expect("raw delivery decodes");
let reader = worker.reader;
let envelope = RawEnvelope {
subject: &delivery.subject,
headers: Some(&delivery.headers),
bytes: delivery.item.0.as_slice(),
sent_at: delivery.sent_at,
handle: delivery.handle,
};
assert!(matches!(
reader.admit_bytes(
&mut worker.scheduler,
&mut worker.tracker,
envelope,
Instant::now(),
),
RawDisposition::Queued { .. }
));
}
async fn run_matcher(bus: MemoryBus, shutdown: Shutdown) {
let mut jobs = bus.source::<SolveRequest<E>>(REQUESTS);
let publisher = bus.publisher::<SolveResult<E>>();
loop {
tokio::select! {
biased;
() = shutdown.triggered() => break,
maybe = jobs.next() => match maybe {
Some(Ok(delivery)) => {
let reply = delivery
.headers
.get(crate::orchestrator::dispatch::REPLY_HEADER)
.expect("requests carry a reply subject")
.as_str()
.to_owned();
let job = delivery.item.resolve(|_, _| None).expect("full request");
let outcome = SolveOutcome::Solved {
diff: MatchedDiff {
revision: 1,
downgraded: false,
layers: Vec::new(),
},
trip: Trip::new(),
converged_through: None,
trip_digest: None,
};
let result = SolveResult::new(&job, outcome, 0);
let _ = publisher
.publish(&reply, &result.msg_id(), HeaderMap::new(), &result)
.await;
let _ = delivery.handle.ack().await;
}
_ => break,
},
}
}
}
async fn wait_until(cond: impl Fn() -> bool) {
for _ in 0..2000 {
if cond() {
return;
}
tokio::time::sleep(Duration::from_millis(2)).await;
}
panic!("condition never held");
}
#[tokio::test(start_paused = true)]
async fn happy_path_three_observations_one_vehicle() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
let mut seqs = Vec::new();
for i in 0..3 {
seqs.push(
publish_raw(
&bus,
partition,
vehicle,
1_775_000_000_000_000 + i * 1_000_000,
)
.await,
);
}
let last_seq = *seqs.last().unwrap();
let worker = build_worker(
config(partition),
Arc::new(catalog(30_000)),
&bus,
&store,
shutdown.clone(),
);
let out_subject = output_subject(u64::from(partition));
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let out_subject = out_subject.clone();
async move {
wait_until(|| bus.published(&out_subject).len() >= 3).await;
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, (), ()) = tokio::join!(
worker.run(),
run_matcher(bus.clone(), shutdown.clone()),
driver
);
let stats = stats.expect("worker ran");
assert_eq!(stats.dispatched, 3, "three jobs dispatched in order");
assert_eq!(stats.committed, 3, "three commits");
assert_eq!(stats.accepted, 3);
assert_eq!(
stats.frontier, last_seq,
"frontier at the last raw sequence"
);
assert_eq!(bus.published(&out_subject).len(), 3, "three outputs");
let (checkpoint, prepared) = store.load(VehicleId(vehicle)).await.unwrap();
assert_eq!(checkpoint.unwrap().revision, Revision(last_seq));
assert!(
prepared.is_none(),
"no prepared record survives an idle vehicle"
);
assert_eq!(bus.acked_count(&raw_subject(u64::from(partition))), 3);
}
#[tokio::test(start_paused = true)]
async fn two_vehicles_interleave_without_blocking() {
let a = 1u64;
let b = same_partition_as(a);
let partition = partition_for(a);
assert_eq!(partition_for(b), partition);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
publish_raw(&bus, partition, a, 1_775_000_000_000_000).await;
publish_raw(&bus, partition, b, 1_775_000_000_000_000).await;
publish_raw(&bus, partition, a, 1_775_000_001_000_000).await;
publish_raw(&bus, partition, b, 1_775_000_001_000_000).await;
let worker = build_worker(
config(partition),
Arc::new(catalog(30_000)),
&bus,
&store,
shutdown.clone(),
);
let out_subject = output_subject(u64::from(partition));
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let out_subject = out_subject.clone();
async move {
wait_until(|| bus.published(&out_subject).len() >= 4).await;
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, (), ()) = tokio::join!(
worker.run(),
run_matcher(bus.clone(), shutdown.clone()),
driver
);
let stats = stats.expect("worker ran");
assert_eq!(stats.committed, 4, "both vehicles fully commit");
assert!(store.load(VehicleId(a)).await.unwrap().0.is_some());
assert!(store.load(VehicleId(b)).await.unwrap().0.is_some());
}
#[tokio::test]
async fn independent_vehicle_commits_overlap_durable_io() {
let a = 1u64;
let b = same_partition_as(a);
let partition = partition_for(a);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
let catalog = Arc::new(catalog(30_000));
publish_raw(&bus, partition, a, 1_775_000_000_000_000).await;
publish_raw(&bus, partition, b, 1_775_000_000_000_000).await;
let output = OverlapPublisher::new(bus.publisher::<CommittedOutput<E>>());
let dispatcher = Dispatcher::new(
bus.publisher::<SolveRequest<E>>(),
DispatchConfig::default(),
REPLY.to_owned(),
);
let committer = Committer::new(
store.clone(),
output.clone(),
CommitConfig {
publish_attempts: 1,
backoff: Duration::ZERO,
..CommitConfig::default()
},
);
let report = clean_report(partition);
let mut worker = PartitionWorker::new(
config(partition),
catalog.clone(),
admission_for(&catalog),
store,
dispatcher,
committer,
bus.source::<RawBytes>(&raw_subject(u64::from(partition))),
bus.source::<SolveResult<E>>(&results_subject(partition)),
&report,
shutdown.clone(),
);
enqueue_next_raw(&mut worker).await;
enqueue_next_raw(&mut worker).await;
worker.pump().await;
let result_publisher = bus.publisher::<SolveResult<E>>();
for (_, _, bytes) in bus.published(REQUESTS) {
let job = full_job(&bytes);
let result = SolveResult::new(
&job,
SolveOutcome::Solved {
diff: MatchedDiff {
revision: 1,
downgraded: false,
layers: Vec::new(),
},
trip: Trip::new(),
converged_through: None,
trip_digest: None,
},
0,
);
result_publisher
.publish(
&results_subject(partition),
&result.msg_id(),
HeaderMap::new(),
&result,
)
.await
.expect("result publishes");
}
for _ in 0..2 {
let delivery = worker
.results
.next()
.await
.expect("result source stays open")
.expect("result decodes");
worker.on_result(delivery).await;
}
assert_eq!(worker.commits.len(), 2);
for _ in 0..2 {
let completion = tokio::time::timeout(Duration::from_secs(1), worker.commits.next())
.await
.unwrap_or_else(|_| {
panic!(
"commit gate stalled: active={}, peak={}, outputs={}",
output.active.load(Ordering::SeqCst),
output.peak.load(Ordering::SeqCst),
bus.published(&output_subject(u64::from(partition))).len(),
)
})
.expect("a commit remains");
worker.finish_commit(completion).await;
}
assert_eq!(worker.stats.committed, 2);
assert_eq!(output.peak.load(Ordering::SeqCst), 2);
}
fn trip_at(origin: Origin) -> Trip<E> {
use routers_network::mock::MockNetworkBuilder;
use routers_transition::Matcher;
use routers_transition::costing::{
CostingStrategies, DefaultEmissionCost, DefaultTransitionCost,
};
use routers_transition::layer::generation::StandardGenerator;
use routers_transition::weigh::AllCompute;
let net = MockNetworkBuilder::new()
.node(1, Point::new(151.2090, -33.8688))
.node(2, Point::new(151.2100, -33.8688))
.edge(1, 2)
.build();
let costing = CostingStrategies::<DefaultEmissionCost, DefaultTransitionCost, E>::default();
let generator = StandardGenerator::new(&net, &costing.emission);
let matcher = Matcher::new(&net, &costing, generator, AllCompute::default(), &());
let mut trip = matcher.begin();
matcher
.push(&mut trip, origin)
.expect("the fixture point anchors");
trip
}
#[tokio::test]
async fn a_sticky_answer_steers_the_next_request_and_a_trip_miss_resends_in_full() {
let vehicle = 1_u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let catalog = Arc::new(catalog(30_000));
let mut worker = build_worker(config(partition), catalog, &bus, &store, Shutdown::new());
let replica = "solve.req.v1.g.g1.r.r1.m.replica";
let group = "solve.req.v1.g.*.r.*.q.*";
let t0 = 1_775_000_000_000_000;
publish_raw_at(&bus, partition, vehicle, t0, point()).await;
enqueue_next_raw(&mut worker).await;
worker.pump().await;
let sent = bus.published(group);
assert_eq!(sent.len(), 1);
let first = SolveRequest::<E>::decode(&sent[0].2)
.unwrap()
.resolve(|_, _| None)
.expect("a first request carries everything");
let result = SolveResult::new(
&first,
SolveOutcome::Solved {
diff: MatchedDiff {
revision: 1,
downgraded: false,
layers: Vec::new(),
},
trip: trip_at(Origin::new(point(), t0)),
converged_through: None,
trip_digest: None,
},
0,
);
let mut headers = HeaderMap::new();
headers.insert(STICKY_HEADER, replica);
bus.publisher::<SolveResult<E>>()
.publish(
&results_subject(partition),
&result.msg_id(),
headers,
&result,
)
.await
.unwrap();
let delivery = worker.results.next().await.unwrap().unwrap();
worker.on_result(delivery).await;
let completion = worker.commits.next().await.expect("the answer commits");
worker.finish_commit(completion).await;
assert_eq!(worker.stats.committed, 1);
publish_raw_at(&bus, partition, vehicle, t0 + 5_000_000, point()).await;
enqueue_next_raw(&mut worker).await;
worker.pump().await;
let steered = bus.published(replica);
assert_eq!(steered.len(), 1, "the resume was steered to the replica");
let cached = SolveRequest::<E>::decode(&steered[0].2).unwrap();
assert!(cached.is_cached());
assert_eq!(bus.published(group).len(), 1, "the group saw nothing new");
let (retained, deadline) = &worker.requests[&VehicleId(vehicle)];
assert!(
!retained.subject.contains(".m."),
"a retry targets the group"
);
assert!(*deadline <= Instant::now() + worker.cfg.sticky_retry);
let miss = SolveResult::<E>::trip_miss(cached.proof.clone(), 0);
bus.publisher::<SolveResult<E>>()
.publish(
&results_subject(partition),
&miss.msg_id(),
HeaderMap::new(),
&miss,
)
.await
.unwrap();
let delivery = worker.results.next().await.unwrap().unwrap();
worker.on_result(delivery).await;
assert!(worker.commits.is_empty(), "a trip miss is never committed");
let resent = bus.published(group);
assert_eq!(resent.len(), 2, "the full request was resent to the group");
let full = SolveRequest::<E>::decode(&resent[1].2).unwrap();
assert!(!full.is_cached());
assert_eq!(full.job_id(), cached.job_id(), "the resend is the same job");
assert!(!worker.sticky.contains_key(&VehicleId(vehicle)));
}
#[tokio::test(start_paused = true)]
async fn overdue_freshness_target_keeps_unanswered_work_active() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
publish_raw(&bus, partition, vehicle, 1_775_000_000_000_000).await;
let worker = build_worker(
config(partition),
Arc::new(catalog(100)),
&bus,
&store,
shutdown.clone(),
);
let out_subject = output_subject(u64::from(partition));
let job_filter = REQUESTS;
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let out_subject = out_subject.clone();
async move {
wait_until(|| !bus.published(job_filter).is_empty()).await;
tokio::time::advance(Duration::from_secs(1)).await;
assert!(
bus.published(&out_subject).is_empty(),
"freshness age must not emit a terminal output"
);
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, ()) = tokio::join!(worker.run(), driver);
let stats = stats.expect("worker ran");
assert_eq!(stats.terminal, 0, "freshness is not a terminal outcome");
assert_eq!(stats.committed, 0);
assert!(bus.published(&out_subject).is_empty());
}
#[tokio::test(start_paused = true)]
async fn duplicate_raw_is_coalesced_and_left_broker_owned() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
publish_raw(&bus, partition, vehicle, 1_775_000_000_000_000).await;
let worker = build_worker(
config(partition),
Arc::new(catalog(30_000)),
&bus,
&store,
shutdown.clone(),
);
let job_filter = REQUESTS;
let raw_sub = raw_subject(u64::from(partition));
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let raw_sub = raw_sub.clone();
async move {
wait_until(|| !bus.published(job_filter).is_empty()).await;
bus.redeliver_unacked(&raw_sub);
wait_until(|| !bus.nak_delays().is_empty()).await;
assert_eq!(bus.acked_count(&raw_sub), 0);
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, ()) = tokio::join!(worker.run(), driver);
let stats = stats.expect("worker ran");
assert_eq!(stats.observed, 2, "the observation was delivered twice");
assert_eq!(stats.coalesced, 1, "the redelivery coalesced");
assert_eq!(stats.dispatched, 1, "only one job was dispatched");
}
#[tokio::test(start_paused = true)]
async fn shutdown_leaves_no_prepared_records_when_idle() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
publish_raw(&bus, partition, vehicle, 1_775_000_000_000_000).await;
publish_raw(&bus, partition, vehicle, 1_775_000_001_000_000).await;
let worker = build_worker(
config(partition),
Arc::new(catalog(30_000)),
&bus,
&store,
shutdown.clone(),
);
let out_subject = output_subject(u64::from(partition));
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let out_subject = out_subject.clone();
async move {
wait_until(|| bus.published(&out_subject).len() >= 2).await;
shutdown.trigger(crate::lifecycle::DrainReason::Signal);
}
};
let (stats, (), ()) = tokio::join!(
worker.run(),
run_matcher(bus.clone(), shutdown.clone()),
driver
);
let stats = stats.expect("worker ran");
assert_eq!(stats.committed, 2);
assert!(
store.snapshot().prepared.is_empty(),
"no prepared records remain"
);
}
async fn seed_checkpoint_bytes(
store: &MemoryCheckpointStore,
vehicle: u64,
partition: u16,
rev: u64,
garbage: Vec<u8>,
) {
let seed_bus = MemoryBus::new();
let obs = ObservationId {
partition,
sequence: rev,
};
let out_subject = output_subject(u64::from(partition));
let output = CommittedOutput::<E>::new(
JobId(u128::from(rev)),
VehicleId(vehicle),
obs,
Revision(rev),
SegmentId(rev),
OutputKind::Terminal {
reason: TerminalReason::Unanchored,
closes_segment: false,
},
);
let entries: Vec<(String, String, Vec<u8>)> = vec![(
out_subject.clone(),
output.msg_id(),
output.encode().unwrap(),
)];
let prepared = PreparedCommit {
output: output.id,
output_subject: out_subject,
output_bytes: postcard::to_allocvec(&entries).unwrap(),
next_checkpoint: garbage,
next_revision: Revision(rev),
next_segment: SegmentId(rev),
expected_base: None,
phase: CommitPhase::Prepared,
raw: obs,
};
store
.prepare(VehicleId(vehicle), partition, prepared.clone())
.await
.unwrap();
let committer = Committer::new(
store.clone(),
seed_bus.publisher::<CommittedOutput<E>>(),
CommitConfig {
publish_attempts: 4,
backoff: Duration::ZERO,
..CommitConfig::default()
},
);
committer
.finish_prepared(VehicleId(vehicle), partition, prepared)
.await
.expect("seed installs the checkpoint bytes");
}
#[tokio::test(start_paused = true)]
async fn state_lost_reset_supersedes_the_stale_checkpoint_without_looping() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
seed_checkpoint_bytes(&store, vehicle, partition, 5, b"not a checkpoint".to_vec()).await;
publish_raw(&bus, partition, vehicle, 1_775_000_000_000_000).await;
let worker = build_worker(
config(partition),
Arc::new(catalog(30_000)),
&bus,
&store,
shutdown.clone(),
);
let out_subject = output_subject(u64::from(partition));
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let out_subject = out_subject.clone();
async move {
wait_until(|| bus.published(&out_subject).len() >= 2).await;
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, (), ()) = tokio::join!(
worker.run(),
run_matcher(bus.clone(), shutdown.clone()),
driver
);
let stats = stats.expect("worker ran");
assert_eq!(stats.committed, 1, "exactly one commit, no loop");
assert_eq!(stats.resets, 1, "the commit opened a fresh segment");
assert_eq!(stats.conflicts, 0, "the CAS targeted the stale revision");
assert_eq!(stats.quarantined_vehicles, 0);
let published = bus.published(&out_subject);
assert_eq!(published.len(), 2, "a reset then a match");
let first = CommittedOutput::<E>::decode(&published[0].2).unwrap();
match first.kind {
OutputKind::Reset { reason, .. } => assert_eq!(reason, ResetReason::StateLost),
other => panic!("expected a StateLost reset first, got {}", other.kind()),
}
let second = CommittedOutput::<E>::decode(&published[1].2).unwrap();
assert!(
matches!(second.kind, OutputKind::Matched { .. }),
"the match follows the reset"
);
let (checkpoint, prepared) = store.load(VehicleId(vehicle)).await.unwrap();
assert!(checkpoint.is_some(), "a fresh checkpoint was promoted");
assert!(prepared.is_none(), "nothing is left staged");
assert_eq!(bus.acked_count(&raw_subject(u64::from(partition))), 1);
}
#[tokio::test(start_paused = true)]
async fn a_poison_prepared_record_quarantines_the_vehicle_instead_of_spinning() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
let prepared = PreparedCommit {
output: OutputId(1),
output_subject: output_subject(u64::from(partition)),
output_bytes: vec![0xff, 0xff, 0xff, 0xff],
next_checkpoint: Vec::new(),
next_revision: Revision(1),
next_segment: SegmentId(1),
expected_base: None,
phase: CommitPhase::Prepared,
raw: ObservationId {
partition,
sequence: 1,
},
};
store
.prepare(VehicleId(vehicle), partition, prepared)
.await
.unwrap();
let admission = admission_for(&catalog(30_000));
let report = RecoveryReport {
partition,
frontier: None,
prepared_found: 1,
prepared_finished: 0,
prepared_failed: vec![VehicleId(vehicle)],
epoch: None,
};
let worker = build_worker_with(
config(partition),
Arc::new(catalog(30_000)),
&bus,
&store,
admission,
&report,
shutdown.clone(),
);
let driver = {
let shutdown = shutdown.clone();
async move {
tokio::time::sleep(Duration::from_millis(200)).await;
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, ()) = tokio::join!(worker.run(), driver);
let stats = stats.expect("worker ran");
assert_eq!(
stats.quarantined_vehicles, 1,
"the poison vehicle was quarantined exactly once"
);
assert_eq!(stats.committed, 0, "nothing could be committed");
assert!(
store.snapshot().prepared.contains_key(&VehicleId(vehicle)),
"the undecodable record is left in place"
);
}
fn unserved_catalog() -> Catalog {
let elsewhere = shard_of(Point::new(0.0, 0.0)).to_string();
assert_ne!(elsewhere, shard_of(point()).to_string());
let toml = format!(
r#"
version = 1
routing_version = 1
[[regions]]
id = "r1"
graph = "g1"
coverage = ["{elsewhere}"]
overlap = []
lanes = 1
resource_class = "c"
replicas = {{ min = 1, max = 1 }}
freshness_budget_ms = 30000
"#,
);
Catalog::parse(&toml).expect("catalog parses")
}
#[tokio::test(start_paused = true)]
async fn an_unserved_point_bypasses_exhausted_admission() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
publish_raw(&bus, partition, vehicle, 1_775_000_000_000_000).await;
let catalog = Arc::new(unserved_catalog());
let admission = Admission::new(
AdmissionConfig {
global_jobs: 0,
region_jobs: 0,
..AdmissionConfig::default()
},
catalog.regions.iter().map(|region| ®ion.id),
);
let admission_probe = admission.clone();
let report = clean_report(partition);
let worker = build_worker_with(
config(partition),
catalog,
&bus,
&store,
admission,
&report,
shutdown.clone(),
);
let out_subject = output_subject(u64::from(partition));
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let out_subject = out_subject.clone();
async move {
wait_until(|| !bus.published(&out_subject).is_empty()).await;
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, ()) = tokio::join!(worker.run(), driver);
let stats = stats.expect("worker ran");
assert_eq!(stats.terminal, 1, "one terminal was committed");
assert_eq!(stats.committed, 1);
assert_eq!(stats.dispatched, 0, "no real job was dispatched");
let published = bus.published(&out_subject);
assert_eq!(published.len(), 1);
let output = CommittedOutput::<E>::decode(&published[0].2).unwrap();
match output.kind {
OutputKind::Terminal { reason, .. } => {
assert_eq!(reason, TerminalReason::UnsupportedCoverage);
}
other => panic!("expected a Terminal output, got {}", other.kind()),
}
assert_eq!(admission_probe.global().jobs, 0, "no job credit leaked");
assert!(
admission_probe.snapshot().iter().all(|r| r.jobs == 0),
"no region credit leaked",
);
}
#[tokio::test(start_paused = true)]
async fn a_state_lost_unserved_point_still_lands_its_terminal_without_looping() {
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
seed_checkpoint_bytes(&store, vehicle, partition, 7, b"not a checkpoint".to_vec()).await;
publish_raw(&bus, partition, vehicle, 1_775_000_000_000_000).await;
let catalog = Arc::new(unserved_catalog());
let admission = admission_for(&catalog);
let admission_probe = admission.clone();
let report = clean_report(partition);
let worker = build_worker_with(
config(partition),
catalog,
&bus,
&store,
admission,
&report,
shutdown.clone(),
);
let out_subject = output_subject(u64::from(partition));
let driver = {
let bus = bus.clone();
let shutdown = shutdown.clone();
let out_subject = out_subject.clone();
async move {
wait_until(|| !bus.published(&out_subject).is_empty()).await;
shutdown.trigger(crate::lifecycle::DrainReason::Operator);
}
};
let (stats, ()) = tokio::join!(worker.run(), driver);
let stats = stats.expect("worker ran");
assert_eq!(stats.terminal, 1, "the terminal landed");
assert_eq!(stats.committed, 1, "exactly one commit, no loop");
assert_eq!(stats.conflicts, 0, "the CAS targeted the stale revision");
assert_eq!(stats.quarantined_vehicles, 0);
let published = bus.published(&out_subject);
assert_eq!(published.len(), 1, "exactly one terminal output");
let output = CommittedOutput::<E>::decode(&published[0].2).unwrap();
match output.kind {
OutputKind::Terminal { reason, .. } => {
assert_eq!(reason, TerminalReason::UnsupportedCoverage);
}
other => panic!("expected a Terminal output, got {}", other.kind()),
}
let (checkpoint, prepared) = store.load(VehicleId(vehicle)).await.unwrap();
assert!(checkpoint.is_some(), "a fresh checkpoint was promoted");
assert!(prepared.is_none(), "nothing is left staged");
assert_eq!(bus.acked_count(&raw_subject(u64::from(partition))), 1);
assert_eq!(admission_probe.global().jobs, 0, "no job credit leaked");
}
#[tokio::test(start_paused = true)]
async fn tick_records_the_oldest_pending_gauge() {
use opentelemetry::metrics::MeterProvider as _;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::error::OTelSdkResult;
use opentelemetry_sdk::metrics::data::{Gauge as GaugeData, ResourceMetrics};
use opentelemetry_sdk::metrics::exporter::PushMetricExporter;
use opentelemetry_sdk::metrics::reader::MetricReader as _;
use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider, Temporality};
#[derive(Debug, Default)]
struct NoopExporter;
impl PushMetricExporter for NoopExporter {
async fn export(&self, _metrics: &mut ResourceMetrics) -> OTelSdkResult {
Ok(())
}
fn force_flush(&self) -> OTelSdkResult {
Ok(())
}
fn shutdown(&self) -> OTelSdkResult {
Ok(())
}
fn temporality(&self) -> Temporality {
Temporality::Cumulative
}
}
fn gauge_value(reader: &PeriodicReader<NoopExporter>, name: &str) -> Option<f64> {
let mut rm = ResourceMetrics {
resource: Resource::builder().build(),
scope_metrics: Vec::new(),
};
reader.collect(&mut rm).expect("collect");
for scope in &rm.scope_metrics {
for metric in &scope.metrics {
if metric.name.as_ref() == name
&& let Some(gauge) = metric.data.as_any().downcast_ref::<GaugeData<f64>>()
{
return gauge.data_points.first().map(|dp| dp.value);
}
}
}
None
}
let vehicle = 1u64;
let partition = partition_for(vehicle);
let bus = MemoryBus::new();
let store = MemoryCheckpointStore::new();
let shutdown = Shutdown::new();
let reader = PeriodicReader::builder(NoopExporter).build();
let provider = SdkMeterProvider::builder()
.with_reader(reader.clone())
.with_resource(Resource::builder().with_service_name("test").build())
.build();
let metrics = Metrics::from_meter(provider.meter("routers_realtime"));
publish_raw(&bus, partition, vehicle, 1_775_000_000_000_000).await;
let mut worker = build_worker(
config(partition),
Arc::new(catalog(3_600_000)),
&bus,
&store,
shutdown.clone(),
)
.with_metrics(metrics);
let delivery = worker
.raw
.next()
.await
.expect("a raw delivery")
.expect("a well-formed delivery");
worker.on_raw(delivery).await;
worker.on_tick().await;
let fresh = gauge_value(&reader, "oldest_pending_age_seconds")
.expect("the tick records the gauge even when the queue just filled");
assert!(fresh < 1.0, "a just-queued head is near zero, got {fresh}");
tokio::time::sleep(Duration::from_secs(5)).await;
worker.on_tick().await;
let aged =
gauge_value(&reader, "oldest_pending_age_seconds").expect("the tick records the gauge");
assert!(
aged >= 5.0,
"the gauge reflects the queued head's age, got {aged}",
);
}
}