use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime};
use tokio::sync::watch;
use tokio::task::JoinHandle;
use super::{Claim, ConsumerId, Delivery, JournalError, JournalStats, PoisonReason, UsageJournal};
use crate::usage::{ObservedRecord, UsageSink};
pub const DRAIN_MARGIN: Duration = Duration::from_secs(1);
const CLOSING_READ: Duration = Duration::from_millis(500);
const PROBE_WRITES: usize = 32;
#[derive(Debug, Clone)]
pub struct WorkerSettings {
pub consumer: ConsumerId,
pub claim_batch: usize,
pub lease: Duration,
pub poll_interval: Duration,
pub maintain_interval: Duration,
}
impl Default for WorkerSettings {
fn default() -> Self {
Self {
consumer: ConsumerId::parse("billing").expect("a static consumer name"),
claim_batch: 256,
lease: Duration::from_secs(30),
poll_interval: Duration::from_millis(250),
maintain_interval: Duration::from_secs(60),
}
}
}
#[derive(Debug, Default, PartialEq, Eq)]
pub struct DrainReport {
pub delivered: u64,
pub failed: u64,
pub undelivered: u64,
pub quarantined: u64,
pub drained: bool,
pub reported: bool,
pub counted: bool,
pub unwaited: bool,
}
impl DrainReport {
pub fn caught_up(&self) -> Option<bool> {
self.counted.then_some(self.drained)
}
fn observe(&mut self, stats: JournalStats) {
self.undelivered = stats.pending + stats.in_flight;
self.quarantined = stats.quarantined;
self.drained = stats.is_drained();
self.counted = true;
}
pub fn log(&self) {
if self.unwaited {
tracing::warn!(
"usage journal delivery was stopped without a drain because the shutdown budget \
was spent; whatever was left is durable and will be delivered after restart"
);
return;
}
if !self.reported {
if self.counted {
tracing::error!(
undelivered = self.undelivered,
quarantined = self.quarantined,
"usage journal worker did not stop within the shutdown bound; the backlog \
below is durable and will be delivered after restart, and this run's \
delivered count is unknown"
);
} else {
tracing::error!(
"usage journal worker did not stop within the shutdown bound and its backlog \
could not be read; the events are durable and will be delivered after restart"
);
}
return;
}
if !self.counted {
tracing::warn!(
delivered = self.delivered,
failed = self.failed,
"usage journal worker stopped but its backlog could not be read; whether \
delivery was caught up is unknown, and anything left behind is durable and \
will be delivered after restart"
);
return;
}
if self.drained {
tracing::info!(
delivered = self.delivered,
quarantined = self.quarantined,
"usage journal drained on shutdown"
);
} else {
tracing::warn!(
delivered = self.delivered,
undelivered = self.undelivered,
quarantined = self.quarantined,
failed = self.failed,
"usage journal was not drained within the shutdown bound; the events are \
durable and will be delivered after restart"
);
}
}
}
pub struct DeliveryWorker {
journal: Arc<dyn UsageJournal>,
sinks: Arc<Vec<Box<dyn UsageSink>>>,
advisory: Arc<Vec<Box<dyn UsageSink>>>,
settings: WorkerSettings,
}
impl DeliveryWorker {
pub fn new(
journal: Arc<dyn UsageJournal>,
sinks: Arc<Vec<Box<dyn UsageSink>>>,
settings: WorkerSettings,
) -> Self {
Self {
journal,
sinks,
advisory: Arc::new(Vec::new()),
settings,
}
}
pub fn also_telling(mut self, advisory: Arc<Vec<Box<dyn UsageSink>>>) -> Self {
self.advisory = advisory;
self
}
pub fn spawn(self) -> WorkerHandle {
let (stop, receiver) = watch::channel(None);
let (journal, consumer) = (Arc::clone(&self.journal), self.settings.consumer.clone());
WorkerHandle {
stop,
task: tokio::spawn(self.run(receiver)),
journal,
consumer,
}
}
async fn run(self, mut stop: watch::Receiver<Option<Duration>>) -> DrainReport {
let mut report = DrainReport {
reported: true,
..DrainReport::default()
};
let mut next_maintain = Instant::now() + self.settings.maintain_interval;
let budget = loop {
if let Some(budget) = self
.pump_until_idle(&mut report, &mut stop, &mut next_maintain)
.await
{
break budget;
}
if !stop.has_changed().unwrap_or(true) {
self.maintain_if_due(&mut next_maintain).await;
}
tokio::select! {
changed = stop.changed() => {
break changed.ok().and_then(|()| *stop.borrow_and_update()).unwrap_or_default();
}
_ = tokio::time::sleep(self.settings.poll_interval) => {}
}
};
let deadline = Instant::now() + budget;
while let Some(remaining) = deadline.checked_duration_since(Instant::now()) {
match tokio::time::timeout(remaining, self.pump(&mut report)).await {
Ok(Ok(0)) | Ok(Err(_)) | Err(_) => break,
Ok(Ok(_)) => {}
}
}
if let Ok(Ok(stats)) =
tokio::time::timeout(CLOSING_READ, self.journal.stats(&self.settings.consumer)).await
{
report.observe(stats);
}
report
}
async fn pump_until_idle(
&self,
report: &mut DrainReport,
stop: &mut watch::Receiver<Option<Duration>>,
next_maintain: &mut Instant,
) -> Option<Duration> {
loop {
if stop.has_changed().unwrap_or(true) {
return None;
}
self.maintain_if_due(next_maintain).await;
let delivered = tokio::select! {
biased;
changed = stop.changed() => {
return Some(
changed.ok().and_then(|()| *stop.borrow_and_update()).unwrap_or_default(),
);
}
delivered = self.pump(report) => delivered,
};
match delivered {
Ok(0) => return None,
Ok(_) => {}
Err(error) => {
tracing::warn!(
journal = self.journal.name(),
error = %error,
"usage journal claim failed; retrying after the poll interval"
);
return None;
}
}
}
}
async fn pump(&self, report: &mut DrainReport) -> Result<usize, JournalError> {
let claimed = self
.journal
.claim(
&self.settings.consumer,
Claim {
max_events: self.settings.claim_batch,
lease: self.settings.lease,
now: SystemTime::now(),
},
)
.await?;
if claimed.is_empty() {
return Ok(0);
}
let redeliveries = claimed
.iter()
.filter(|delivery| delivery.id.is_redelivery())
.count();
if redeliveries > 0 {
crate::telemetry::metrics::record_usage_journal_deliveries(
self.journal.name(),
self.settings.consumer.as_str(),
"redelivered",
redeliveries as u64,
);
}
let outcome = self.deliver(&claimed).await;
if outcome.landed.len() < claimed.len() {
report.failed += 1;
crate::telemetry::metrics::record_usage_journal_deliveries(
self.journal.name(),
self.settings.consumer.as_str(),
"failed",
(claimed.len() - outcome.landed.len()) as u64,
);
let attempts = self.journal.capacity().max_delivery_attempts;
for index in &outcome.refused {
let delivery = &claimed[*index];
if delivery.id.attempt < attempts {
continue;
}
match self
.journal
.quarantine(&delivery.id, PoisonReason::Rejected)
.await
{
Ok(()) => crate::telemetry::metrics::record_usage_journal_quarantined(
self.journal.name(),
self.settings.consumer.as_str(),
PoisonReason::Rejected.as_str(),
),
Err(error) => tracing::warn!(
delivery = %delivery.id,
error = %error,
"usage event exhausted its delivery attempts but could not be quarantined"
),
}
}
for (index, delivery) in claimed.iter().enumerate() {
if outcome.landed.contains(&index) || outcome.refused.contains(&index) {
continue;
}
if let Err(error) = self.journal.relinquish(&delivery.id).await {
tracing::warn!(
delivery = %delivery.id,
error = %error,
"usage event's delivery attempt could not be returned after an \
unattributable refusal"
);
}
}
}
if outcome.landed.is_empty() {
return Ok(0);
}
let landed: Vec<super::DeliveryId> = outcome
.landed
.iter()
.map(|index| claimed[*index].id.clone())
.collect();
let mut acknowledged = 0;
for (delivery, verdict) in landed.iter().zip(self.journal.ack_all(&landed).await) {
match verdict {
Ok(()) => acknowledged += 1,
Err(error) => tracing::warn!(
delivery = %delivery,
error = %error,
"usage event was written but not acknowledged; it will be redelivered"
),
}
}
report.delivered += acknowledged as u64;
crate::telemetry::metrics::record_usage_journal_deliveries(
self.journal.name(),
self.settings.consumer.as_str(),
"acknowledged",
acknowledged as u64,
);
Ok(acknowledged)
}
async fn deliver(&self, claimed: &[Delivery]) -> Delivered {
let outcome = self.deliver_durably(claimed).await;
self.count_written(outcome.landed.len());
self.tell_advisory(claimed, &outcome.landed).await;
outcome
}
fn count_written(&self, landed: usize) {
if landed == 0 {
return;
}
for sink in self.sinks.iter() {
crate::telemetry::metrics::record_usage_written(sink.name(), landed as u64);
}
}
async fn deliver_durably(&self, claimed: &[Delivery]) -> Delivered {
let mut outcome = Delivered::default();
if self.write(claimed).await {
outcome.landed.extend(0..claimed.len());
return outcome;
}
if claimed.len() == 1 {
return outcome;
}
let mut orphans: Vec<usize> = Vec::new();
let mut refusals = 1;
let mut level = vec![halve(0..claimed.len())];
while !level.is_empty() {
let mut next = Vec::new();
for (left, right) in level {
for range in [left, right] {
if self.write(&claimed[range.clone()]).await {
outcome.landed.extend(range);
continue;
}
refusals += 1;
if range.len() == 1 {
orphans.push(range.start);
} else {
next.push(halve(range));
}
if outcome.landed.is_empty() && refusals >= PROBE_WRITES {
return Delivered::default();
}
}
}
level = next;
}
if outcome.landed.is_empty() {
return outcome;
}
outcome.refused = orphans;
outcome
}
async fn tell_advisory(&self, claimed: &[Delivery], landed: &[usize]) {
if self.advisory.is_empty() || landed.is_empty() {
return;
}
let batch: Vec<ObservedRecord> = landed
.iter()
.map(|&index| claimed[index].event.observed())
.collect();
for sink in self.advisory.iter() {
match sink.record_batch(&batch).await {
Ok(()) => {
crate::telemetry::metrics::record_usage_written(sink.name(), batch.len() as u64)
}
Err(error) => tracing::warn!(
sink = sink.name(),
records = batch.len(),
error = %error,
"usage telemetry export failed; the events are delivered and stay delivered"
),
}
}
}
async fn write(&self, claimed: &[Delivery]) -> bool {
let batch: Vec<ObservedRecord> = claimed
.iter()
.map(|delivery| delivery.event.observed())
.collect();
for sink in self.sinks.iter() {
if let Err(error) = sink.record_batch(&batch).await {
tracing::warn!(
sink = sink.name(),
records = batch.len(),
error = %error,
"usage journal delivery failed; the events stay journaled"
);
return false;
}
}
true
}
async fn publish_stats(&self) {
match self.journal.stats(&self.settings.consumer).await {
Ok(stats) => crate::telemetry::metrics::record_usage_journal_stats(
self.journal.name(),
self.settings.consumer.as_str(),
&stats,
),
Err(error) => tracing::warn!(
journal = self.journal.name(),
error = %error,
"usage journal stats could not be read"
),
}
}
async fn maintain_if_due(&self, next_maintain: &mut Instant) {
if Instant::now() < *next_maintain {
return;
}
self.maintain().await;
self.publish_stats().await;
*next_maintain = Instant::now() + self.settings.maintain_interval;
}
async fn maintain(&self) {
match self.journal.maintain(SystemTime::now()).await {
Ok(pruned) if pruned > 0 => {
tracing::debug!(
journal = self.journal.name(),
pruned,
"usage journal pruned"
)
}
Ok(_) => {}
Err(error) => tracing::warn!(
journal = self.journal.name(),
error = %error,
"usage journal retention pass failed"
),
}
self.report_other_consumers().await;
}
async fn report_other_consumers(&self) {
match self
.journal
.consumers_besides(&self.settings.consumer)
.await
{
Ok(others) if !others.is_empty() => tracing::warn!(
journal = self.journal.name(),
consumer = %self.settings.consumer,
others = others.join(","),
"the usage journal is holding delivery state for consumers this deployment is \
not running, and retention waits on every one of them; if they are retired \
names, delete their rows, or the outbox grows to its limit and refuses appends"
),
Ok(_) => {}
Err(error) => tracing::debug!(
journal = self.journal.name(),
error = %error,
"usage journal consumers could not be listed"
),
}
}
}
#[derive(Default)]
struct Delivered {
landed: Vec<usize>,
refused: Vec<usize>,
}
fn halve(range: std::ops::Range<usize>) -> (std::ops::Range<usize>, std::ops::Range<usize>) {
let mid = range.start + range.len() / 2;
(range.start..mid, mid..range.end)
}
pub struct WorkerHandle {
stop: watch::Sender<Option<Duration>>,
task: JoinHandle<DrainReport>,
journal: Arc<dyn UsageJournal>,
consumer: ConsumerId,
}
impl WorkerHandle {
pub fn abandon(self) -> DrainReport {
let _ = self.stop.send(Some(Duration::ZERO));
DrainReport {
unwaited: true,
..DrainReport::default()
}
}
pub async fn drain(self, budget: Duration) -> DrainReport {
let Self {
stop,
task,
journal,
consumer,
} = self;
let _ = stop.send(Some(budget));
match tokio::time::timeout(budget + DRAIN_MARGIN, task).await {
Ok(Ok(report)) => report,
Ok(Err(error)) => {
tracing::error!(error = %error, "usage journal worker panicked");
Self::unreported(&journal, &consumer).await
}
Err(_) => {
tracing::error!("usage journal worker did not stop inside its bound");
Self::unreported(&journal, &consumer).await
}
}
}
async fn unreported(journal: &Arc<dyn UsageJournal>, consumer: &ConsumerId) -> DrainReport {
let mut report = DrainReport::default();
match tokio::time::timeout(CLOSING_READ, journal.stats(consumer)).await {
Ok(Ok(stats)) => report.observe(stats),
Err(_) => tracing::error!(
journal = journal.name(),
"usage journal backlog could not be read inside its bound after the drain was \
abandoned"
),
Ok(Err(error)) => tracing::error!(
journal = journal.name(),
error = %error,
"usage journal backlog could not be read after the drain was abandoned"
),
}
report.drained = false;
report
}
}
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use async_trait::async_trait;
use super::super::oracle::InMemoryUsageJournal;
use super::super::tests::{consumer, event_for};
use super::super::{Appended, Capacity, CapacityPolicy, DeliveryId, PoisonReason, UsageEvent};
use super::*;
use crate::usage::{SinkFailure, UsageRecord};
#[derive(Default)]
struct Recorder {
written: Mutex<Vec<String>>,
refuse: bool,
poison: Vec<String>,
seen: Mutex<Vec<Vec<String>>>,
}
impl Recorder {
fn written(&self) -> Vec<String> {
self.written.lock().expect("not poisoned").clone()
}
fn batches(&self) -> Vec<Vec<String>> {
self.seen.lock().expect("not poisoned").clone()
}
}
#[async_trait]
impl UsageSink for Recorder {
fn name(&self) -> &'static str {
"recorder"
}
async fn record(&self, record: &UsageRecord) {
self.written
.lock()
.expect("not poisoned")
.push(record.request_id.clone());
}
async fn record_batch(&self, batch: &[ObservedRecord]) -> Result<(), SinkFailure> {
self.seen.lock().expect("not poisoned").push(
batch
.iter()
.map(|observed| observed.record.request_id.clone())
.collect(),
);
if self.refuse {
return Err(SinkFailure::new("the destination is refusing writes"));
}
if batch
.iter()
.any(|observed| self.poison.contains(&observed.record.request_id))
{
return Err(SinkFailure::new("the destination refuses this row"));
}
for observed in batch {
self.record(&observed.record).await;
}
Ok(())
}
}
struct LosesAcks(InMemoryUsageJournal);
#[async_trait]
impl UsageJournal for LosesAcks {
fn name(&self) -> &'static str {
"loses-acks"
}
fn capacity(&self) -> Capacity {
self.0.capacity()
}
fn mode(&self) -> super::super::DeliveryMode {
self.0.mode()
}
async fn append(&self, event: &UsageEvent) -> Result<Appended, JournalError> {
self.0.append(event).await
}
async fn claim(
&self,
consumer: &ConsumerId,
claim: Claim,
) -> Result<Vec<Delivery>, JournalError> {
self.0.claim(consumer, claim).await
}
async fn ack(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
Err(JournalError::Backend(format!(
"the acknowledgement of {delivery} never reached the journal"
)))
}
async fn quarantine(
&self,
delivery: &DeliveryId,
reason: PoisonReason,
) -> Result<(), JournalError> {
self.0.quarantine(delivery, reason).await
}
async fn relinquish(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
self.0.relinquish(delivery).await
}
async fn stats(&self, consumer: &ConsumerId) -> Result<JournalStats, JournalError> {
self.0.stats(consumer).await
}
}
struct SlowStats(InMemoryUsageJournal, Duration);
#[async_trait]
impl UsageJournal for SlowStats {
fn name(&self) -> &'static str {
"slow-stats"
}
fn capacity(&self) -> Capacity {
self.0.capacity()
}
fn mode(&self) -> super::super::DeliveryMode {
self.0.mode()
}
async fn append(&self, event: &UsageEvent) -> Result<Appended, JournalError> {
self.0.append(event).await
}
async fn claim(
&self,
consumer: &ConsumerId,
claim: Claim,
) -> Result<Vec<Delivery>, JournalError> {
self.0.claim(consumer, claim).await
}
async fn ack(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
self.0.ack(delivery).await
}
async fn quarantine(
&self,
delivery: &DeliveryId,
reason: PoisonReason,
) -> Result<(), JournalError> {
self.0.quarantine(delivery, reason).await
}
async fn relinquish(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
self.0.relinquish(delivery).await
}
async fn stats(&self, consumer: &ConsumerId) -> Result<JournalStats, JournalError> {
tokio::time::sleep(self.1).await;
self.0.stats(consumer).await
}
}
#[tokio::test]
async fn a_slow_closing_backlog_read_does_not_make_a_stopped_worker_look_abandoned() {
let journal = Arc::new(SlowStats(
InMemoryUsageJournal::new(),
Duration::from_secs(5),
));
let sinks: Vec<Box<dyn UsageSink>> = vec![];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(sinks),
settings(Duration::from_secs(30)),
)
.spawn();
let started = Instant::now();
let report = handle.drain(Duration::from_millis(50)).await;
assert!(
report.reported,
"a worker abandoned for its closing read: {report:?}"
);
assert!(
started.elapsed() < Duration::from_secs(2),
"shutdown waited on the journal's own timeout: {:?}",
started.elapsed()
);
assert!(
!report.counted,
"the backlog read did not finish, so it is unknown rather than zero: {report:?}"
);
assert_eq!(
report.caught_up(),
None,
"an unread backlog cannot be reported as an incomplete drain: {report:?}"
);
}
struct SlowMaintain(InMemoryUsageJournal, Arc<AtomicBool>, Duration);
#[async_trait]
impl UsageJournal for SlowMaintain {
fn name(&self) -> &'static str {
"slow-maintain"
}
fn capacity(&self) -> Capacity {
self.0.capacity()
}
fn mode(&self) -> super::super::DeliveryMode {
self.0.mode()
}
async fn append(&self, event: &UsageEvent) -> Result<Appended, JournalError> {
self.0.append(event).await
}
async fn claim(
&self,
consumer: &ConsumerId,
claim: Claim,
) -> Result<Vec<Delivery>, JournalError> {
self.0.claim(consumer, claim).await
}
async fn ack(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
self.0.ack(delivery).await
}
async fn quarantine(
&self,
delivery: &DeliveryId,
reason: PoisonReason,
) -> Result<(), JournalError> {
self.0.quarantine(delivery, reason).await
}
async fn relinquish(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
self.0.relinquish(delivery).await
}
async fn stats(&self, consumer: &ConsumerId) -> Result<JournalStats, JournalError> {
self.0.stats(consumer).await
}
async fn maintain(&self, now: SystemTime) -> Result<u64, JournalError> {
if self.1.load(Ordering::SeqCst) {
tokio::time::sleep(self.2).await;
}
self.0.maintain(now).await
}
}
struct Announcing(watch::Sender<bool>, Duration);
#[async_trait]
impl UsageSink for Announcing {
fn name(&self) -> &'static str {
"announcing"
}
async fn record(&self, _record: &UsageRecord) {
let _ = self.0.send(true);
tokio::time::sleep(self.1).await;
}
async fn record_batch(&self, _batch: &[ObservedRecord]) -> Result<(), SinkFailure> {
let _ = self.0.send(true);
tokio::time::sleep(self.1).await;
Ok(())
}
}
#[tokio::test]
async fn housekeeping_due_at_shutdown_does_not_cost_the_worker_its_report() {
let slow = Arc::new(AtomicBool::new(false));
let journal = Arc::new(SlowMaintain(
InMemoryUsageJournal::new(),
Arc::clone(&slow),
Duration::from_secs(5),
));
let (writing, mut written) = watch::channel(false);
let sinks: Vec<Box<dyn UsageSink>> =
vec![Box::new(Announcing(writing, Duration::from_millis(100)))];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(sinks),
WorkerSettings {
maintain_interval: Duration::ZERO,
..settings(Duration::from_secs(30))
},
)
.spawn();
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
written.changed().await.expect("the worker started writing");
slow.store(true, Ordering::SeqCst);
let started = Instant::now();
let report = handle.drain(Duration::from_millis(50)).await;
assert!(
report.reported,
"the worker was abandoned inside a housekeeping pass it should not have started: \
{report:?}"
);
assert!(
started.elapsed() < Duration::from_secs(2),
"shutdown waited on housekeeping: {:?}",
started.elapsed()
);
}
#[tokio::test]
async fn a_single_slow_delivery_pass_cannot_overrun_the_shutdown_bound() {
let journal = Arc::new(InMemoryUsageJournal::new());
let sinks: Vec<Box<dyn UsageSink>> = vec![Box::new(Slow(Duration::from_secs(5)))];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(sinks),
WorkerSettings {
poll_interval: Duration::from_secs(30),
..settings(Duration::from_secs(30))
},
)
.spawn();
tokio::time::sleep(Duration::from_millis(20)).await;
for _ in 0..8 {
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
}
let started = Instant::now();
let report = handle.drain(Duration::from_millis(50)).await;
assert!(
report.reported,
"a pass slower than the bound left the worker abandoned: {report:?}"
);
assert!(
started.elapsed() < Duration::from_secs(2),
"shutdown waited on the destination: {:?}",
started.elapsed()
);
assert_eq!(report.undelivered, 8, "{report:?}");
assert!(!report.drained, "{report:?}");
}
#[tokio::test]
async fn a_pass_already_in_flight_when_stop_arrives_does_not_spend_the_bound() {
let journal = Arc::new(InMemoryUsageJournal::new());
let sinks: Vec<Box<dyn UsageSink>> = vec![Box::new(Slow(Duration::from_secs(5)))];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(sinks),
settings(Duration::from_secs(30)),
)
.spawn();
for _ in 0..8 {
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
}
tokio::time::sleep(Duration::from_millis(50)).await;
let started = Instant::now();
let report = handle.drain(Duration::from_millis(50)).await;
assert!(
report.reported,
"a worker that stopped correctly mid-pass was abandoned: {report:?}"
);
assert!(
started.elapsed() < Duration::from_secs(2),
"shutdown waited for the in-flight destination write: {:?}",
started.elapsed()
);
assert_eq!(report.undelivered, 8, "{report:?}");
assert!(!report.drained, "{report:?}");
}
fn settings(lease: Duration) -> WorkerSettings {
WorkerSettings {
consumer: consumer("billing"),
claim_batch: 8,
lease,
poll_interval: Duration::from_millis(5),
maintain_interval: Duration::from_millis(20),
}
}
fn worker(
journal: Arc<dyn UsageJournal>,
sink: Arc<Recorder>,
settings: WorkerSettings,
) -> WorkerHandle {
let sinks: Vec<Box<dyn UsageSink>> = vec![Box::new(SharedSink(Arc::clone(&sink)))];
DeliveryWorker::new(journal, Arc::new(sinks), settings).spawn()
}
struct SharedSink(Arc<Recorder>);
#[async_trait]
impl UsageSink for SharedSink {
fn name(&self) -> &'static str {
self.0.name()
}
async fn record(&self, record: &UsageRecord) {
self.0.record(record).await;
}
async fn record_batch(&self, batch: &[ObservedRecord]) -> Result<(), SinkFailure> {
self.0.record_batch(batch).await
}
}
async fn eventually(
sink: &Recorder,
what: &str,
predicate: impl Fn(&[String]) -> bool,
) -> Vec<String> {
let deadline = Instant::now() + Duration::from_secs(5);
loop {
let written = sink.written();
if predicate(&written) {
return written;
}
assert!(
Instant::now() < deadline,
"{what} did not happen; the sink saw {written:?}"
);
tokio::time::sleep(Duration::from_millis(5)).await;
}
}
#[tokio::test]
async fn an_appended_event_is_delivered_and_acknowledged() {
let journal = Arc::new(InMemoryUsageJournal::new());
let sink = Arc::new(Recorder::default());
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_secs(30)),
);
let event = event_for("GW_INBOUND_ACME_KEY");
journal.append(&event).await.expect("append");
let written = eventually(&sink, "the event reached the sink", |written| {
!written.is_empty()
})
.await;
assert_eq!(written, vec![event.id().to_string()]);
let report = handle.drain(Duration::from_secs(5)).await;
assert_eq!(report.delivered, 1, "{report:?}");
assert_eq!(report.undelivered, 0, "{report:?}");
assert!(report.drained, "{report:?}");
}
#[tokio::test]
async fn a_destination_that_refuses_the_batch_leaves_the_event_journaled() {
let journal = Arc::new(InMemoryUsageJournal::new());
let sink = Arc::new(Recorder {
refuse: true,
..Recorder::default()
});
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_secs(30)),
);
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
tokio::time::sleep(Duration::from_millis(50)).await;
let report = handle.drain(Duration::from_millis(50)).await;
assert!(sink.written().is_empty(), "nothing was written");
assert_eq!(report.delivered, 0, "{report:?}");
assert_eq!(report.undelivered, 1, "{report:?}");
assert!(report.failed > 0, "{report:?}");
assert!(!report.drained, "{report:?}");
assert_eq!(journal.stored_events(), 1);
}
#[tokio::test]
async fn a_destination_wide_outage_quarantines_nothing_however_long_it_lasts() {
let journal = Arc::new(InMemoryUsageJournal::with_capacity(Capacity {
max_events: 8,
max_delivery_attempts: 1,
retain_acknowledged: Duration::from_secs(60),
policy: CapacityPolicy::Refuse,
}));
let sink = Arc::new(Recorder {
refuse: true,
..Recorder::default()
});
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_millis(5)),
);
for subject in ["one", "two", "three", "four"] {
journal.append(&event_for(subject)).await.expect("append");
}
tokio::time::sleep(Duration::from_millis(200)).await;
let report = handle.drain(Duration::from_millis(50)).await;
assert!(report.failed > 1, "the batch was retried: {report:?}");
assert_eq!(report.quarantined, 0, "{report:?}");
assert_eq!(report.undelivered, 4, "{report:?}");
assert_eq!(report.delivered, 0, "{report:?}");
}
#[tokio::test]
async fn one_refused_event_is_isolated_and_its_siblings_are_delivered() {
let journal = Arc::new(InMemoryUsageJournal::with_capacity(Capacity {
max_events: 8,
max_delivery_attempts: 1,
retain_acknowledged: Duration::from_secs(60),
policy: CapacityPolicy::Refuse,
}));
let poison = event_for("poison");
let sink = Arc::new(Recorder {
poison: vec![poison.record().request_id.clone()],
..Recorder::default()
});
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_millis(5)),
);
journal.append(&poison).await.expect("append");
for subject in ["one", "two", "three"] {
journal.append(&event_for(subject)).await.expect("append");
}
let deadline = Instant::now() + Duration::from_secs(5);
loop {
let stats = journal.stats(&consumer("billing")).await.expect("stats");
if stats.quarantined == 1 && stats.pending == 0 && stats.in_flight == 0 {
break;
}
assert!(
Instant::now() < deadline,
"the refused event was not isolated: {stats:?}"
);
tokio::time::sleep(Duration::from_millis(5)).await;
}
let written = sink.written();
assert!(
!written.contains(&poison.record().request_id),
"the refused event was never written: {written:?}"
);
let report = handle.drain(Duration::from_millis(50)).await;
assert_eq!(report.quarantined, 1, "{report:?}");
assert_eq!(report.undelivered, 0, "{report:?}");
}
#[tokio::test]
async fn a_destination_that_cannot_answer_is_told_but_never_acknowledged_on() {
let journal = Arc::new(InMemoryUsageJournal::with_capacity(Capacity {
max_events: 8,
max_delivery_attempts: 1,
retain_acknowledged: Duration::from_secs(60),
policy: CapacityPolicy::Refuse,
}));
let poison = event_for("poison");
let storing = Arc::new(Recorder {
poison: vec![poison.record().request_id.clone()],
..Recorder::default()
});
let exported = Arc::new(Recorder {
refuse: true,
..Recorder::default()
});
let acknowledged: Vec<Box<dyn UsageSink>> =
vec![Box::new(SharedSink(Arc::clone(&storing)))];
let advisory: Vec<Box<dyn UsageSink>> = vec![Box::new(SharedSink(Arc::clone(&exported)))];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(acknowledged),
settings(Duration::from_millis(5)),
)
.also_telling(Arc::new(advisory))
.spawn();
journal.append(&poison).await.expect("append");
for subject in ["one", "two", "three"] {
journal.append(&event_for(subject)).await.expect("append");
}
eventually(&storing, "the healthy events", |written| written.len() == 3).await;
let report = handle.drain(Duration::from_millis(200)).await;
assert_eq!(report.delivered, 3, "{report:?}");
assert_eq!(report.quarantined, 1, "{report:?}");
assert_eq!(report.undelivered, 0, "{report:?}");
assert!(
!exported
.batches()
.iter()
.any(|batch| batch.contains(&poison.record().request_id)),
"an event no destination stored was exported: {:?}",
exported.batches()
);
for batch in exported.batches() {
assert!(
batch.len() <= 3 && !batch.is_empty(),
"the export saw a bisection probe rather than the delivered set: {batch:?}"
);
}
}
#[tokio::test]
async fn two_refused_events_in_one_batch_are_both_isolated() {
let journal = Arc::new(InMemoryUsageJournal::with_capacity(Capacity {
max_events: 8,
max_delivery_attempts: 1,
retain_acknowledged: Duration::from_secs(60),
policy: CapacityPolicy::Refuse,
}));
let (first, second) = (event_for("poison-one"), event_for("poison-two"));
let (healthy, other) = (event_for("one"), event_for("two"));
let sink = Arc::new(Recorder {
poison: vec![
first.record().request_id.clone(),
second.record().request_id.clone(),
],
..Recorder::default()
});
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_millis(5)),
);
for event in [&first, &healthy, &second, &other] {
journal.append(event).await.expect("append");
}
let deadline = Instant::now() + Duration::from_secs(5);
loop {
let stats = journal.stats(&consumer("billing")).await.expect("stats");
if stats.quarantined == 2 && stats.pending == 0 && stats.in_flight == 0 {
break;
}
assert!(
Instant::now() < deadline,
"the refused events were not isolated: {stats:?}"
);
tokio::time::sleep(Duration::from_millis(5)).await;
}
let written = sink.written();
for delivered in [&healthy, &other] {
assert!(
written.contains(&delivered.record().request_id),
"a sibling of the refused events was delivered: {written:?}"
);
}
let report = handle.drain(Duration::from_millis(50)).await;
assert_eq!(report.quarantined, 2, "{report:?}");
assert_eq!(report.undelivered, 0, "{report:?}");
}
struct NeverCatchesUp {
inner: InMemoryUsageJournal,
maintained: AtomicUsize,
}
#[async_trait]
impl UsageJournal for NeverCatchesUp {
fn name(&self) -> &'static str {
"never-catches-up"
}
fn capacity(&self) -> Capacity {
self.inner.capacity()
}
fn mode(&self) -> super::super::DeliveryMode {
self.inner.mode()
}
async fn append(&self, event: &UsageEvent) -> Result<Appended, JournalError> {
self.inner.append(event).await
}
async fn claim(
&self,
consumer: &ConsumerId,
claim: Claim,
) -> Result<Vec<Delivery>, JournalError> {
tokio::time::sleep(Duration::from_millis(1)).await;
self.inner.append(&event_for("arriving")).await?;
self.inner.claim(consumer, claim).await
}
async fn ack(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
self.inner.ack(delivery).await
}
async fn quarantine(
&self,
delivery: &DeliveryId,
reason: PoisonReason,
) -> Result<(), JournalError> {
self.inner.quarantine(delivery, reason).await
}
async fn relinquish(&self, delivery: &DeliveryId) -> Result<(), JournalError> {
self.inner.relinquish(delivery).await
}
async fn stats(&self, consumer: &ConsumerId) -> Result<JournalStats, JournalError> {
self.inner.stats(consumer).await
}
async fn maintain(&self, now: SystemTime) -> Result<u64, JournalError> {
self.maintained.fetch_add(1, Ordering::Relaxed);
self.inner.maintain(now).await
}
}
#[tokio::test]
async fn housekeeping_runs_on_a_worker_that_never_catches_up() {
let journal = Arc::new(NeverCatchesUp {
inner: InMemoryUsageJournal::new(),
maintained: AtomicUsize::new(0),
});
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(Recorder::default()),
settings(Duration::from_secs(30)),
);
let deadline = Instant::now() + Duration::from_secs(5);
while journal.maintained.load(Ordering::Relaxed) == 0 {
assert!(
Instant::now() < deadline,
"a worker that never went idle never ran its maintenance tick"
);
tokio::time::sleep(Duration::from_millis(5)).await;
}
handle.drain(Duration::from_millis(50)).await;
}
#[tokio::test]
async fn a_write_whose_acknowledgement_was_lost_is_delivered_again() {
let journal = Arc::new(LosesAcks(InMemoryUsageJournal::new()));
let sink = Arc::new(Recorder::default());
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_millis(5)),
);
let event = event_for("GW_INBOUND_ACME_KEY");
journal.append(&event).await.expect("append");
let written = eventually(&sink, "the event was delivered twice", |written| {
written.len() >= 2
})
.await;
assert!(
written.iter().all(|id| *id == event.id().to_string()),
"{written:?}"
);
handle.drain(Duration::from_millis(50)).await;
}
#[tokio::test]
async fn one_callers_events_are_delivered_in_append_order() {
let journal = Arc::new(InMemoryUsageJournal::new());
let sink = Arc::new(Recorder::default());
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_secs(30)),
);
let mut expected = Vec::new();
for _ in 0..4 {
let event = event_for("GW_INBOUND_ACME_KEY");
expected.push(event.id().to_string());
journal.append(&event).await.expect("append");
}
let written = eventually(&sink, "every event reached the sink", move |written| {
written.len() >= 4
})
.await;
assert_eq!(written, expected);
handle.drain(Duration::from_millis(50)).await;
}
#[tokio::test]
async fn a_drain_that_runs_out_of_budget_reports_the_backlog_rather_than_waiting() {
let journal = Arc::new(InMemoryUsageJournal::new());
let sink = Arc::new(Recorder {
refuse: true,
..Recorder::default()
});
let handle = worker(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::clone(&sink),
settings(Duration::from_secs(30)),
);
for _ in 0..3 {
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
}
let started = Instant::now();
let report = handle.drain(Duration::from_millis(50)).await;
assert!(
started.elapsed() < Duration::from_secs(3),
"the drain waited past its bound: {:?}",
started.elapsed()
);
assert_eq!(report.undelivered, 3, "{report:?}");
assert!(!report.drained, "{report:?}");
assert!(
report.reported,
"the worker reported for itself: {report:?}"
);
}
#[tokio::test]
async fn a_drain_that_had_to_abandon_the_worker_reports_the_backlog_rather_than_zeros() {
let slow = Arc::new(AtomicBool::new(true));
let journal = Arc::new(SlowMaintain(
InMemoryUsageJournal::new(),
Arc::clone(&slow),
Duration::from_secs(60),
));
let sinks: Vec<Box<dyn UsageSink>> = vec![Box::new(Recorder::default())];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(sinks),
WorkerSettings {
maintain_interval: Duration::ZERO,
..settings(Duration::from_secs(30))
},
)
.spawn();
for _ in 0..3 {
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
}
tokio::time::sleep(Duration::from_millis(50)).await;
let report = handle.drain(Duration::from_millis(10)).await;
assert!(!report.reported, "the worker never reported: {report:?}");
assert!(report.counted, "{report:?}");
assert_eq!(report.undelivered, 3, "{report:?}");
assert!(!report.drained, "{report:?}");
}
#[tokio::test]
async fn a_worker_can_be_stopped_without_spending_a_drain_margin_on_it() {
let journal = Arc::new(InMemoryUsageJournal::new());
let sinks: Vec<Box<dyn UsageSink>> = vec![Box::new(Slow(Duration::from_secs(30)))];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(sinks),
settings(Duration::from_secs(30)),
)
.spawn();
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
let started = Instant::now();
let report = handle.abandon();
assert!(
started.elapsed() < DRAIN_MARGIN,
"abandoning waited on the worker: {:?}",
started.elapsed()
);
assert!(
!report.reported && !report.counted,
"a report nobody produced must claim nothing: {report:?}"
);
assert!(!report.drained, "{report:?}");
assert!(
report.unwaited,
"a deliberate stop must not read as an overrun: {report:?}"
);
}
struct Slow(Duration);
#[async_trait]
impl UsageSink for Slow {
fn name(&self) -> &'static str {
"slow"
}
async fn record(&self, _record: &UsageRecord) {
tokio::time::sleep(self.0).await;
}
async fn record_batch(&self, _batch: &[ObservedRecord]) -> Result<(), SinkFailure> {
tokio::time::sleep(self.0).await;
Ok(())
}
}
#[tokio::test]
async fn a_worker_with_a_long_backlog_stops_at_its_bound_instead_of_being_abandoned() {
let journal = Arc::new(InMemoryUsageJournal::with_capacity(Capacity {
max_events: 4096,
max_delivery_attempts: 8,
retain_acknowledged: Duration::from_secs(60),
policy: CapacityPolicy::Refuse,
}));
let sinks: Vec<Box<dyn UsageSink>> = vec![Box::new(Slow(Duration::from_millis(30)))];
let handle = DeliveryWorker::new(
Arc::clone(&journal) as Arc<dyn UsageJournal>,
Arc::new(sinks),
settings(Duration::from_secs(30)),
)
.spawn();
for _ in 0..400 {
journal
.append(&event_for("GW_INBOUND_ACME_KEY"))
.await
.expect("append");
}
let report = handle.drain(Duration::from_millis(50)).await;
assert!(
report.reported,
"the worker stopped at its bound and reported: {report:?}"
);
assert!(report.undelivered > 0, "{report:?}");
assert!(!report.drained, "{report:?}");
}
}