Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
f65712e
perf(ltx): overlap verified compaction I/O and measure sustained writes
forhappy Oct 3, 2026
f9995af
Merge branch 'main' into codex/write-publication-performance
forhappy Oct 3, 2026
b1a141f
test(ltx): verify cancellation retains parallel flush resources
forhappy Oct 3, 2026
325b5d6
docs(perf): retain sustained write results and unresolved bottleneck
forhappy Oct 3, 2026
1b73491
perf: measure write admission and provider costs
forhappy Oct 4, 2026
51e4446
perf(ltx): avoid reserving dirty capacity behind recovery waits
forhappy Oct 4, 2026
ec19f8c
test(app): snapshot object counters and samples at one boundary
forhappy Oct 4, 2026
ddf7461
docs(perf): retain measured admission and provider bottlenecks
forhappy Oct 4, 2026
8a66af7
test(ltx): attribute warm writes and compaction storage reads
forhappy Oct 4, 2026
caa8a01
docs(perf): retain isolated recovery admission measurements
forhappy Oct 4, 2026
16ac5d8
perf(runtime): group queued native mutations behind one durable root
forhappy Oct 4, 2026
1e21c17
fix(perf): recognize byte-identical instrumented baselines
forhappy Oct 4, 2026
0a31048
perf(axum): add sustained queued-write comparison profile
forhappy Oct 4, 2026
0647bce
docs(perf): retain paired native group write gains and regressions
forhappy Oct 4, 2026
8614596
fix(runtime): release completed request slots before reply delivery
forhappy Oct 4, 2026
4a98b37
test(runtime): keep reply handoff gate compatible with stable Rust
forhappy Oct 4, 2026
0b4b4c2
perf(ltx): refill compaction transfers before earlier sources finish
forhappy Oct 4, 2026
3d21bb5
test(peer-http): account for grouped publication ranges
forhappy Oct 4, 2026
98ba666
docs(perf): retain audited C64 paired write measurements
forhappy Oct 4, 2026
05a77ee
docs(perf): centralize concise reports and tracked metrics
forhappy Oct 4, 2026
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
74 changes: 66 additions & 8 deletions .github/workflows/axum-performance.yml
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
name: Paired Axum read performance
name: Paired Axum HTTP performance

on:
workflow_call:
Expand All @@ -9,12 +9,21 @@ on:
read_seconds:
type: number
default: 60
workload:
type: string
default: reads
write_seconds:
type: number
default: 120
write_concurrency:
type: number
default: 16

permissions:
contents: read

jobs:
steady-reads:
steady-http:
runs-on: ubuntu-24.04
timeout-minutes: 70
services:
Expand All @@ -41,11 +50,41 @@ jobs:
- name: Build both release servers before measuring
env:
BASELINE_REF: ${{ inputs.baseline_ref }}
WORKLOAD: ${{ inputs.workload }}
run: |
set -euo pipefail
mkdir -p "$RUNNER_TEMP/axum-evidence"
baseline_commit=$(git rev-parse --verify --end-of-options "$BASELINE_REF^{commit}")
git worktree add --detach "$RUNNER_TEMP/axum-baseline-source" "$baseline_commit"
if [ "$WORKLOAD" = writes ]; then
# Backport only phase observations and the example observer wiring.
# Retain the exact diff so the baseline's instrumentation is auditable.
patch="$GITHUB_WORKSPACE/scripts/axum-write-telemetry.patch"
# The candidate may reorganize admission; this fixed fixture only
# backports observations, never candidate scheduling changes.
if git -C "$RUNNER_TEMP/axum-baseline-source" apply --check "$patch"; then
git -C "$RUNNER_TEMP/axum-baseline-source" apply "$patch"
elif git -C "$RUNNER_TEMP/axum-baseline-source" apply --reverse --check "$patch"; then
:
else
# An instrumented baseline may have reorganized these files.
# Accept it only when every patch target already matches the
# candidate byte for byte; never copy candidate policy into it.
git apply --numstat "$patch" > "$RUNNER_TEMP/axum-evidence/telemetry-targets.tsv"
test -s "$RUNNER_TEMP/axum-evidence/telemetry-targets.tsv"
while IFS=$'\t' read -r added removed target; do
cmp "$GITHUB_WORKSPACE/$target" "$RUNNER_TEMP/axum-baseline-source/$target"
done < "$RUNNER_TEMP/axum-evidence/telemetry-targets.tsv"
fi
cp crates/cellule-axum/examples/sql_metrics/mod.rs \
"$RUNNER_TEMP/axum-baseline-source/crates/cellule-axum/examples/sql_metrics/mod.rs"
cp "$patch" "$RUNNER_TEMP/axum-evidence/telemetry-backport.patch"
git -C "$RUNNER_TEMP/axum-baseline-source" diff \
> "$RUNNER_TEMP/axum-evidence/baseline-instrumentation.diff"
sha256sum "$RUNNER_TEMP/axum-evidence/"*.patch \
"$RUNNER_TEMP/axum-evidence/baseline-instrumentation.diff" \
> "$RUNNER_TEMP/axum-evidence/instrumentation.sha256"
fi
printf '%s\n' "$baseline_commit" > "$RUNNER_TEMP/axum-evidence/baseline-revision.txt"
git rev-parse HEAD > "$RUNNER_TEMP/axum-evidence/candidate-revision.txt"
uname -a > "$RUNNER_TEMP/axum-evidence/host.txt"
Expand All @@ -57,9 +96,14 @@ jobs:
--manifest-path "$RUNNER_TEMP/axum-baseline-source/Cargo.toml"
CARGO_TARGET_DIR="$RUNNER_TEMP/axum-candidate-target" \
cargo build --release -p cellule-axum --examples --locked
- name: Measure paired steady reads and verify every Cell
- name: Measure paired steady HTTP and verify every Cell
env:
READ_SECONDS: ${{ inputs.read_seconds }}
WRITE_SECONDS: ${{ inputs.write_seconds }}
WRITE_CONCURRENCY: ${{ inputs.write_concurrency }}
WORKLOAD: ${{ inputs.workload }}
TOKIO_WORKER_THREADS: "4"
RUSTFS_CONTAINER_ID: ${{ job.services.rustfs.id }}
AWS_ACCESS_KEY_ID: cellule-ci
AWS_SECRET_ACCESS_KEY: cellule-integration-fixture
AWS_DEFAULT_REGION: us-east-1
Expand All @@ -69,20 +113,34 @@ jobs:
CELLULE_TEST_PREFIX: paired-${{ github.run_id }}-${{ github.run_attempt }}
run: |
set -euo pipefail
python3 -c 'import os; assert 30 <= int(os.environ["READ_SECONDS"]) <= 120'
concurrency=16
case "$WORKLOAD" in
reads)
python3 -c 'import os; assert 30 <= int(os.environ["READ_SECONDS"]) <= 120'
load_args=(--read-driver "$RUNNER_TEMP/axum-candidate-target/release/examples/http_load"
--read-warmup-seconds 5 --read-seconds "$READ_SECONDS")
;;
writes)
python3 -c 'import os; assert 60 <= int(os.environ["WRITE_SECONDS"]) <= 180'
python3 -c 'import os; assert int(os.environ["WRITE_CONCURRENCY"]) in (16, 64)'
concurrency="$WRITE_CONCURRENCY"
load_args=(--write-warmup-seconds 5 --write-seconds "$WRITE_SECONDS"
--provider-container "$RUSTFS_CONTAINER_ID")
;;
*) exit 1 ;;
esac
aws --endpoint-url "$CELLULE_TEST_ENDPOINT" s3api create-bucket --bucket "$CELLULE_TEST_BUCKET"
python3 scripts/bench-axum-rustfs.py \
--binary "$RUNNER_TEMP/axum-candidate-target/release/examples/sql" \
--baseline-binary "$RUNNER_TEMP/axum-baseline-target/release/examples/sql" \
--read-driver "$RUNNER_TEMP/axum-candidate-target/release/examples/http_load" \
--output "$RUNNER_TEMP/axum-evidence/results" \
--repeats 3 --cells 1 4 16 --workers 4 --concurrency 16 \
--repeats 3 --cells 1 4 16 --workers 4 --concurrency "$concurrency" \
--warmup 16 --writes 32 --reads 64 \
--read-warmup-seconds 5 --read-seconds "$READ_SECONDS"
"${load_args[@]}"
- uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4
if: always()
with:
name: axum-steady-${{ github.run_id }}-${{ github.run_attempt }}
name: axum-steady-${{ inputs.workload }}-${{ github.run_id }}-${{ github.run_attempt }}
path: ${{ runner.temp }}/axum-evidence/
if-no-files-found: error
retention-days: 7
30 changes: 24 additions & 6 deletions .github/workflows/write-capacity.yml
Original file line number Diff line number Diff line change
Expand Up @@ -15,18 +15,27 @@ on:
workflow_dispatch:
inputs:
mode:
description: Qualification or paired Axum read measurements
description: Qualification or paired Axum HTTP measurements
type: choice
default: qualification
options: [qualification, axum-reads]
options: [qualification, axum-reads, axum-writes]
baseline_ref:
description: Baseline with the multicell Axum SQL example
description: Baseline with the multicell Axum SQL example and telemetry
type: string
default: a277e5282badab55ceb58433fdddf0dee4dc8542
default: 4c627d3fdc95fbb86558088927aec72f51d0046c
read_seconds:
description: Seconds per steady-state read phase (30 to 120)
type: number
default: 60
write_seconds:
description: Seconds per steady-state write phase (60 to 180)
type: number
default: 120
write_concurrency:
description: Write clients; 64 additionally exercises per-Cell queueing
type: choice
default: "16"
options: ["16", "64"]

permissions:
contents: read
Expand All @@ -37,7 +46,7 @@ concurrency:

jobs:
object-proof:
if: github.event_name != 'workflow_dispatch' || inputs.mode != 'axum-reads'
if: github.event_name != 'workflow_dispatch' || inputs.mode == 'qualification'
runs-on: ubuntu-24.04
timeout-minutes: 100
steps:
Expand Down Expand Up @@ -104,8 +113,17 @@ jobs:
baseline_ref: ${{ inputs.baseline_ref }}
read_seconds: ${{ fromJSON(format('{0}', inputs.read_seconds || 60)) }}

axum-writes:
if: github.event_name == 'workflow_dispatch' && inputs.mode == 'axum-writes'
uses: ./.github/workflows/axum-performance.yml
with:
baseline_ref: ${{ inputs.baseline_ref }}
workload: writes
write_seconds: ${{ fromJSON(format('{0}', inputs.write_seconds || 120)) }}
write_concurrency: ${{ fromJSON(format('{0}', inputs.write_concurrency || 16)) }}

follower-proof:
if: github.event_name != 'workflow_dispatch' || inputs.mode != 'axum-reads'
if: github.event_name != 'workflow_dispatch' || inputs.mode == 'qualification'
runs-on: ubuntu-24.04
timeout-minutes: 100
steps:
Expand Down
6 changes: 6 additions & 0 deletions crates/cellule-app/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,12 @@
Declare a stable Cell topology, compile native Rust modules, and expose typed
application handles. The host owns runtime lifecycle and network wiring.

Each SQL Cell owns its SQLite database, request ledger, and capture/recovery
lineage. Cells can share worker threads and resource budgets within a host,
and run on different nodes under separate fenced ownership. The application
builder declares Cell types, partitions, schemas, and operations; the host
handles placement, activation, routing, and recovery.

```mermaid
flowchart LR
Modules[Native modules] --> Builder[ApplicationBuilder]
Expand Down
66 changes: 62 additions & 4 deletions crates/cellule-app/tests/entities/process/observation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,20 +22,39 @@ pub(super) struct StorageCounters {
samples: Mutex<Vec<(i64, StorageObservation)>>,
}

struct StorageSnapshot {
outcomes: [[u64; 9]; 11],
samples: Vec<(i64, StorageObservation)>,
}

impl StorageCounters {
fn snapshot(&self) -> StorageSnapshot {
// Peer listeners remain available while other owners drain. Capture
// counters and samples at one boundary even if a peer read completes.
let samples = self.samples.lock().unwrap();
let outcomes = std::array::from_fn(|operation| {
std::array::from_fn(|outcome| self.outcomes[operation][outcome].load(Ordering::Relaxed))
});
let samples = samples.clone();
StorageSnapshot { outcomes, samples }
}
}

impl StorageObserver for StorageCounters {
fn started(&self, _: StorageOperation) {
self.started.fetch_add(1, Ordering::Relaxed);
}

fn finished(&self, observation: StorageObservation) {
let mut samples = self.samples.lock().unwrap();
self.outcomes[observation.operation.index()][observation.outcome.index()]
.fetch_add(1, Ordering::Relaxed);
self.bytes_read
.fetch_add(observation.bytes_read, Ordering::Relaxed);
self.bytes_written
.fetch_add(observation.bytes_written, Ordering::Relaxed);
self.finished.fetch_add(1, Ordering::Relaxed);
self.samples.lock().unwrap().push((now_ms(), observation));
samples.push((now_ms(), observation));
}
}

Expand Down Expand Up @@ -138,11 +157,11 @@ impl NodeObservations {
}

pub(super) fn finish(&mut self, storage: &StorageCounters, durability: &DurabilityRecorder) {
let storage = storage.snapshot();
writeln!(self.objects, "operation\toutcome\tcount").unwrap();
for operation in StorageOperation::ALL {
for outcome in StorageOutcome::ALL {
let count =
storage.outcomes[operation.index()][outcome.index()].load(Ordering::Relaxed);
let count = storage.outcomes[operation.index()][outcome.index()];
writeln!(
self.objects,
"{}\t{}\t{count}",
Expand All @@ -157,7 +176,7 @@ impl NodeObservations {
"at_ms\toperation\toutcome\tduration_us\tbytes_read\tbytes_written"
)
.unwrap();
for (at_ms, observation) in storage.samples.lock().unwrap().iter() {
for (at_ms, observation) in &storage.samples {
writeln!(
self.object_operations,
"{at_ms}\t{}\t{}\t{}\t{}\t{}",
Expand Down Expand Up @@ -286,6 +305,45 @@ impl NodeObservations {
}
}

#[test]
fn object_counter_snapshots_match_samples_during_peer_completions() {
let storage = StorageCounters::default();
let barrier = std::sync::Barrier::new(5);
std::thread::scope(|scope| {
for _ in 0..4 {
scope.spawn(|| {
barrier.wait();
for _ in 0..2_000 {
storage.finished(StorageObservation {
operation: StorageOperation::Get,
outcome: StorageOutcome::Success,
duration: Duration::from_micros(1),
bytes_read: 1,
bytes_written: 0,
});
std::thread::yield_now();
}
});
}
barrier.wait();
for _ in 0..100 {
let snapshot = storage.snapshot();
let mut samples = [[0_u64; 9]; 11];
for (_, observation) in snapshot.samples {
samples[observation.operation.index()][observation.outcome.index()] += 1;
}
assert_eq!(snapshot.outcomes, samples);
std::thread::yield_now();
}
});
let final_snapshot = storage.snapshot();
assert_eq!(final_snapshot.samples.len(), 8_000);
assert_eq!(
final_snapshot.outcomes[StorageOperation::Get.index()][StorageOutcome::Success.index()],
8_000
);
}

fn disk_bytes(root: &Path) -> u64 {
std::fs::read_dir(root)
.unwrap()
Expand Down
58 changes: 58 additions & 0 deletions crates/cellule-axum/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -560,3 +560,61 @@ gh workflow run write-capacity.yml --ref YOUR_BRANCH \
-f baseline_ref=a277e5282badab55ceb58433fdddf0dee4dc8542 \
-f read_seconds=60
```

For sustained writes, use timed POST admission after five seconds of write
warmup. Every issued identity and response is retained, and every acknowledged
row and exact retry is checked live and after a fresh-file cold restore:

```sh
python3 "$task_dir/source/scripts/bench-axum-rustfs.py" \
--binary "$CARGO_TARGET_DIR/release/examples/sql" \
--output "$task_dir/write-results" \
--repeats 3 --cells 1 4 16 --workers 4 --concurrency 16 \
--warmup 16 --writes 32 --reads 64 \
--write-warmup-seconds 5 --write-seconds 120

gh workflow run write-capacity.yml --ref YOUR_BRANCH \
-f mode=axum-writes -f baseline_ref=BASELINE_COMMIT -f write_seconds=120
```

The workflow defaults to sixteen write clients. Run an additional comparison
with `-f write_concurrency=64` to exercise per-Cell queues at sixteen Cells.
Keep its results separate from the sixteen-client series: concurrency changes
the workload. Both profiles retain three alternating pairs per Cell count,
the same resource ceilings, and every durability and cold-recovery check.

Timed writes replace the fixed `--writes` count. Clients stop admitting POSTs
at the deadline and wait for every in-flight request; throughput includes that
drain time. JSON response-body throughput excludes headers and transport.
The closed-loop Python write driver is intended for storage-bound commands.
Its per-client ledgers retain uncertain transport outcomes without retrying
them; any error fails verification. Each phase retains at most 100,000 responses.
The example reports lifetime publication, compaction, worker and queue timing
histograms after drain, including warmup and correctness checks. Root preparation
separates dirty-memory admission from admitted work; its total includes both.
Dirty and recovery semaphore waits are also measured across replica operations.
These phases overlap with preparation and compaction totals; do not add them.
It also records effective host capacities and fixed provider-operation duration,
outcome, and byte counters. Provider counts include startup, verification, and
maintenance, so they are not steady-window write request counts. Histogram
upper bounds have 100 µs resolution through two seconds; overflow percentiles
are unknown. The paired write workflow instruments both servers identically,
retaining its observation-only baseline patch and complete baseline diff in the
artifact. The patch changes no admission ceilings or publication barriers.
The write workflow also records dedicated RustFS cgroup CPU counters before
and after each load window, including warmup. For a local dedicated container,
pass `--provider-container CONTAINER_ID`; this requires cgroup v2 and Docker.

The [sustained-write report](performance/2026-10-03-rustfs-steady-writes.md)
retains all paired results and durability evidence. The initial compaction
comparison improves compaction time but does not establish consistent gains
in HTTP write throughput or tail latency. The current-main confirmation also
retains mixed results; neither comparison establishes the optimization goal.

The later [group-commit report](performance/2026-10-04-rustfs-grouped-paired-writes.md)
records sustained throughput and latency gains at one and four Cells, together
with every sixteen-Cell regression and the failed 64-client baseline warmup.

```sh
cargo test -p cellule-axum --all-targets --locked
```
11 changes: 8 additions & 3 deletions crates/cellule-axum/examples/sql.rs
Original file line number Diff line number Diff line change
Expand Up @@ -340,7 +340,13 @@ async fn main() -> ExampleResult<()> {
.map(|shard| CellTarget::new(tenant, application_id, ORDERS, &partition_for_shard(shard)))
.collect::<cellule_runtime::Result<_>>()?;
let (store, prefix) = example_storage().await?;
let layout = CellStorageLayout::new(store, prefix, *application_id.as_bytes());
let host = Host::default().with_local_disk_budget(DiskBudget::new(1 << 30));
let query_metrics = Arc::new(sql_metrics::QueryMetrics::new(&host));
let layout = CellStorageLayout::new(
store.with_storage_observer(query_metrics.clone()),
prefix,
*application_id.as_bytes(),
);
let registry = application.registry();
let code = registry
.module_code(Orders::NAME)
Expand All @@ -357,9 +363,8 @@ async fn main() -> ExampleResult<()> {
SqlWorkerPool::new(usize::try_from(workers)?, usize::try_from(MAX_CELLS)?)?,
16 * 1024 * 1024,
session,
Host::default().with_local_disk_budget(DiskBudget::new(1 << 30)),
host,
)?;
let query_metrics = Arc::new(sql_metrics::QueryMetrics::default());
let result: ExampleResult<()> = async {
runtime.install_telemetry(query_metrics.clone())?;
let mut handles = Vec::with_capacity(targets.len());
Expand Down
Loading
Loading