Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
/target
.DS_Store
node_modules/
__pycache__/
*.pyc
63 changes: 63 additions & 0 deletions crates/cellule-axum/examples/capacity/arrivals/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
//! Emit the complete offered schedule, including offers whose wakeup is late.
use super::{Job, Kind};
use std::time::Duration;
use tokio::{sync::mpsc, time::Instant};

#[cfg(test)]
mod tests;

#[derive(Default)]
pub(super) struct Arrivals {
pub offered: [u64; 2],
pub dropped: [u64; 2],
pub warmup_dropped: [u64; 2],
}

pub(super) async fn produce(
sender: mpsc::Sender<Job>,
start: Instant,
warm_end: Instant,
end: Instant,
rates: [u64; 2],
) -> Arrivals {
let mut counts = Arrivals::default();
let mut indexes = [0_u64; 2];
loop {
let write_due = (indexes[0] * 1_000_000_000)
.checked_div(rates[0])
.map_or(end, |nanos| start + Duration::from_nanos(nanos));
let read_due = (indexes[1] * 1_000_000_000)
.checked_div(rates[1])
.map_or(end, |nanos| start + Duration::from_nanos(nanos));
let kind = usize::from(read_due < write_due);
let due = if kind == 0 { write_due } else { read_due };
if due >= end {
break;
}
// Wakeup delay cannot erase an offer scheduled inside the window.
// Its original due time still controls latency and completion gates;
// a full queue records a drop rather than hiding work as unissued.
tokio::time::sleep_until(due).await;
let measured = due >= warm_end;
counts.offered[kind] += u64::from(measured);
let job = Job {
kind: if kind == 0 { Kind::Write } else { Kind::Read },
index: indexes[kind],
due,
measured,
};
match sender.try_send(job) {
Ok(()) => {}
Err(mpsc::error::TrySendError::Full(_)) => {
if measured {
counts.dropped[kind] += 1;
} else {
counts.warmup_dropped[kind] += 1;
}
}
Err(mpsc::error::TrySendError::Closed(_)) => break,
}
indexes[kind] += 1;
}
counts
}
77 changes: 77 additions & 0 deletions crates/cellule-axum/examples/capacity/arrivals/tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
use super::*;

#[tokio::test]
async fn delayed_producer_keeps_every_original_due_time_and_window_offer() {
let start = Instant::now() - Duration::from_secs(5);
let warm_end = start + Duration::from_secs(1);
let end = warm_end + Duration::from_secs(1);
let (sender, mut receiver) = mpsc::channel(20);
let arrivals = produce(sender, start, warm_end, end, [2, 3]).await;
assert_eq!(arrivals.offered, [2, 3]);
assert_eq!(arrivals.dropped, [0, 0]);
assert_eq!(arrivals.warmup_dropped, [0, 0]);
let mut counts = [0; 2];
let mut measured = [0; 2];
while let Some(job) = receiver.recv().await {
let kind = match job.kind {
Kind::Write => 0,
Kind::Read => 1,
};
let rate = [2, 3][kind];
assert_eq!(job.index, counts[kind]);
assert_eq!(
job.due,
start + Duration::from_nanos(job.index * 1_000_000_000 / rate)
);
assert!(job.due < end);
assert_eq!(job.measured, job.due >= warm_end);
counts[kind] += 1;
measured[kind] += u64::from(job.measured);
}
assert_eq!(counts, [4, 6]);
assert_eq!(measured, arrivals.offered);
}

#[tokio::test]
async fn late_full_queue_records_drops_instead_of_omitting_offers() {
let start = Instant::now() - Duration::from_secs(5);
let warm_end = start + Duration::from_secs(1);
let end = warm_end + Duration::from_secs(1);
let (sender, mut receiver) = mpsc::channel(1);
let arrivals = produce(sender, start, warm_end, end, [0, 3]).await;
assert_eq!(arrivals.offered, [0, 3]);
assert_eq!(arrivals.dropped, [0, 3]);
assert_eq!(arrivals.warmup_dropped, [0, 2]);
assert!(!receiver.recv().await.unwrap().measured);
assert!(receiver.recv().await.is_none());
}

#[tokio::test]
async fn immediate_recovery_window_needs_no_new_warmup() {
let start = Instant::now() - Duration::from_secs(5);
let end = start + Duration::from_secs(1);
let (sender, mut receiver) = mpsc::channel(4);
let arrivals = produce(sender, start, start, end, [2, 0]).await;
assert_eq!(arrivals.offered, [2, 0]);
assert_eq!(arrivals.warmup_dropped, [0, 0]);
while let Some(job) = receiver.recv().await {
assert!(job.measured);
}
}

#[test]
fn recovery_metadata_and_zero_warmup_are_supported_by_the_real_config_decoder() {
let config = serde_json::json!({"address":"127.0.0.1:8080", "cells":1000,
"concurrency":128, "queue_capacity":128, "write_rate":10, "read_rate":0,
"warmup_seconds":0, "seconds":30, "evidence_directory":"/tmp/audit",
"phase":"recovery"});
let decoded: super::super::Config = serde_json::from_value(config.clone()).unwrap();
decoded.validate().unwrap();
let mut warmed = config.clone();
warmed["warmup_seconds"] = serde_json::json!(30);
let decoded: super::super::Config = serde_json::from_value(warmed).unwrap();
assert!(decoded.validate().is_err());
let mut unknown = config;
unknown["phase"] = serde_json::json!("hide_steady_point");
assert!(serde_json::from_value::<super::super::Config>(unknown).is_err());
}
Loading
Loading