Skip to content
Draft
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
6 changes: 6 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

8 changes: 7 additions & 1 deletion crates/cellule-axum/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -24,13 +24,19 @@ utoipa-axum = { version = "0.3", optional = true }
[dev-dependencies]
axum = { workspace = true, features = ["http1", "tokio"] }
blake3.workspace = true
bytes.workspace = true
cellule-peer-http.workspace = true
cellule-ltx = { workspace = true, features = ["replica"] }
cellule-store.workspace = true
object_store.workspace = true
reqwest.workspace = true
ed25519-dalek.workspace = true
futures-util.workspace = true
prost.workspace = true
reqwest = { workspace = true, features = ["json"] }
rusqlite.workspace = true
tempfile.workspace = true
tokio = { workspace = true, features = ["macros", "net", "rt-multi-thread", "signal"] }
tokio-util = { workspace = true, features = ["rt"] }
tower = { version = "0.5", features = ["util"] }

[[example]]
Expand Down
14 changes: 12 additions & 2 deletions crates/cellule-axum/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -445,10 +445,10 @@ The same example can use real S3 objects. Set `CELLULE_TEST_ENDPOINT`,
`AWS_ACCESS_KEY_ID`, and `AWS_SECRET_ACCESS_KEY`. `AWS_DEFAULT_REGION`
defaults to `us-east-1`; `AWS_SESSION_TOKEN` is optional.
`CELLULE_AXUM_BIND` selects a loopback listener (default `127.0.0.1:3000`).
`CELLULE_AXUM_CELLS` selects 1–16 active SQL Cells, and
`CELLULE_AXUM_CELLS` selects 1–2,000 active SQL Cells, and
`CELLULE_AXUM_WORKERS` selects 1–16 SQL workers; both default to one.
Order IDs route to shard `id mod active_cells` (Euclidean remainder).
The application declares 16 fixed shards and activates the selected prefix
The application declares 2,048 fixed shards and activates the selected prefix
of that topology. Each active Cell has its own database, writer, publication
root, and request ledger; handlers use `CellClient::local_many` through the
same typed application handle and `Cellule` extractor.
Expand All @@ -462,6 +462,16 @@ each performance point. An active owner is refused;
this tutorial demonstrates graceful restart, not failed-owner takeover.
Keep the binary unchanged so its application code digest matches the catalog.

The optional follower fixture uses this same SQL service and adapter with two
separate processes, pinned mTLS, signed directory enrollment, and fsynced
`FollowerStore` logs. Its application-owned wire protocol leaves canonical
Cellule commit frames unchanged. See the [node capacity runner](performance/node-capacity.md)
for setup, proof-source checks, resource evidence, and remaining qualification
requirements. HTTP endpoints and transport wiring remain application-owned.
If the follow-up read fails after a command commits, the example retains the
committed receipt in its published-failure response. Resolve the original
mutation before issuing a replacement.

The [performance runner](../../scripts/bench-axum-rustfs.py) measures real
HTTP POST and GET requests, using the adapter, application handles, SQL
runtime, LTX publication, and RustFS. It verifies receipt sequences, every
Expand Down
135 changes: 135 additions & 0 deletions crates/cellule-axum/examples/capacity/metrics.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
//! Scheduled latency includes load-generator queueing; HTTP latency starts at dispatch.
use std::time::Duration;
const BUCKET_US: u64 = 100;
const BUCKETS: usize = 100_002;

pub(super) struct Metrics {
pub warmup_attempts: u64,
pub warmup_errors: u64,
pub errors: u64,
attempts: u64,
successes: u64,
completed_in_window: u64,
bytes: u64,
per_cell_successes: Vec<u64>,
scheduled: Vec<u32>,
request: Vec<u32>,
max_scheduled_ms: f64,
max_request_ms: f64,
}

fn observe(histogram: &mut [u32], latency: Duration) {
let bucket = latency.as_micros().div_ceil(u128::from(BUCKET_US));
let index = usize::try_from(bucket)
.unwrap_or(usize::MAX)
.min(BUCKETS - 1);
histogram[index] += 1;
}

fn percentile(histogram: &[u32], attempts: u64, percent: u64) -> Option<f64> {
if attempts == 0 {
return None;
}
let rank = (attempts * percent).div_ceil(100);
let mut count = 0_u64;
for (index, bucket) in histogram.iter().enumerate() {
count += u64::from(*bucket);
if count >= rank {
return (index < BUCKETS - 1).then_some(index as f64 * BUCKET_US as f64 / 1_000.0);
}
}
None
}

impl Metrics {
pub fn new(cells: usize) -> Self {
Self {
warmup_attempts: 0,
warmup_errors: 0,
errors: 0,
attempts: 0,
successes: 0,
completed_in_window: 0,
bytes: 0,
per_cell_successes: vec![0; cells],
scheduled: vec![0; BUCKETS],
request: vec![0; BUCKETS],
max_scheduled_ms: 0.0,
max_request_ms: 0.0,
}
}

pub fn observe(
&mut self,
shard: usize,
success: bool,
in_window: bool,
scheduled: Duration,
request: Duration,
bytes: usize,
) {
self.attempts += 1;
self.errors += u64::from(!success);
self.successes += u64::from(success);
self.completed_in_window += u64::from(success && in_window);
if success {
self.per_cell_successes[shard] += 1;
self.bytes += bytes as u64;
}
observe(&mut self.scheduled, scheduled);
observe(&mut self.request, request);
self.max_scheduled_ms = self.max_scheduled_ms.max(scheduled.as_secs_f64() * 1_000.0);
self.max_request_ms = self.max_request_ms.max(request.as_secs_f64() * 1_000.0);
}

pub fn merge(&mut self, other: Self) {
self.warmup_attempts += other.warmup_attempts;
self.warmup_errors += other.warmup_errors;
self.attempts += other.attempts;
self.successes += other.successes;
self.errors += other.errors;
self.completed_in_window += other.completed_in_window;
self.bytes += other.bytes;
for (a, b) in self
.per_cell_successes
.iter_mut()
.zip(other.per_cell_successes)
{
*a += b;
}
for (a, b) in self.scheduled.iter_mut().zip(other.scheduled) {
*a += b;
}
for (a, b) in self.request.iter_mut().zip(other.request) {
*a += b;
}
self.max_scheduled_ms = self.max_scheduled_ms.max(other.max_scheduled_ms);
self.max_request_ms = self.max_request_ms.max(other.max_request_ms);
}

pub fn summary(
&self,
seconds: u64,
rate: u64,
offered: u64,
dropped: u64,
warmup_dropped: u64,
) -> serde_json::Value {
serde_json::json!({
"target_requests_per_second": rate, "planned_offers": rate*seconds,
"generated_offers": offered, "producer_unissued": (rate*seconds).saturating_sub(offered),
"queue_dropped": dropped, "warmup_queue_dropped": warmup_dropped,
"warmup_attempts": self.warmup_attempts, "warmup_errors": self.warmup_errors,
"attempts": self.attempts, "successes_including_drain": self.successes, "errors": self.errors,
"successes_in_window": self.completed_in_window, "successful_requests_per_second": self.completed_in_window as f64 / seconds as f64,
"successes_after_deadline": self.successes - self.completed_in_window,
"payload_bytes_including_drain": self.bytes, "per_cell_successes_including_drain": self.per_cell_successes,
"scheduled_latency_ms_all_attempts": {"p50": percentile(&self.scheduled,self.attempts,50), "p95": percentile(&self.scheduled,self.attempts,95), "p99": percentile(&self.scheduled,self.attempts,99), "max":self.max_scheduled_ms},
"request_latency_ms_all_attempts": {"p50": percentile(&self.request,self.attempts,50), "p95": percentile(&self.request,self.attempts,95), "p99": percentile(&self.request,self.attempts,99), "max":self.max_request_ms},
"histogram_resolution_us": BUCKET_US, "scheduled_histogram_overflow": self.scheduled[BUCKETS-1], "request_histogram_overflow": self.request[BUCKETS-1],
})
}
}

#[cfg(test)]
mod tests;
39 changes: 39 additions & 0 deletions crates/cellule-axum/examples/capacity/metrics/tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
use super::*;

#[test]
fn delayed_completions_do_not_inflate_window_tps_and_queue_latency_is_visible() {
let mut metrics = Metrics::new(1);
metrics.observe(
0,
true,
true,
Duration::from_millis(12),
Duration::from_millis(2),
100,
);
metrics.observe(
0,
true,
false,
Duration::from_millis(300),
Duration::from_millis(3),
100,
);
metrics.observe(
0,
false,
true,
Duration::from_secs(11),
Duration::from_millis(4),
0,
);
let summary = metrics.summary(1, 5, 4, 1, 0);
assert_eq!(summary["successful_requests_per_second"], 1.0);
assert_eq!(summary["successes_including_drain"], 2);
assert_eq!(summary["successes_after_deadline"], 1);
assert_eq!(summary["errors"], 1);
assert_eq!(summary["producer_unissued"], 1);
assert_eq!(summary["scheduled_histogram_overflow"], 1);
assert!(summary["scheduled_latency_ms_all_attempts"]["p99"].is_null());
assert_eq!(summary["request_latency_ms_all_attempts"]["p99"], 4.0);
}
Loading