Files
smarm/examples/causal_pipeline.rs
T
Claude (sandbox) 2668f4018f feat(causal): native causal profiling behind smarm-causal (RFC 007 v1)
causal_site! scoped site guards per actor slot, progress! throughput
points, and a Coz-style virtual-speedup engine hooked into
maybe_preempt's amortized cold block: target-site samples grow a global
delay ledger; bystanders spin-absorb their debt at the next causal
check, with timeslice extension so injected delay is not charged
against the slice.

Resume credit (Coz's blocked-thread rule) is gated on a causal_parked
slot bit set only by a real park: crediting on every resume made any
yield-cadence actor delay-immune and every experiment inert (found
live on a 24-core run — dead-flat deltas across all sites).

Report normalization uses a measured TSC frequency (~50ms calibration
on first use) instead of the crate-wide 3 GHz assumption, which
uniformly inflated impact numbers on a 3.7 GHz box. impact_pct() is
the machine-readable form of the summary for programmatic checks.

examples/causal_pipeline.rs burns fixed *work* (calibrated LCG loop),
not fixed wall time — a timed busy-wait absorbs injected delay into
its own budget and reads as a no-op. Self-checking: exits nonzero if
causal separation fails; skips the verdict below 4 cores. Validated
on a 24-core box: reserve (true bottleneck) +29.3%@25/+83.5%@50;
serialize and background-compaction ~0%.

Known v1 gaps (jar): timer-heap deadlines unshifted, no-check!/no-alloc
actors undelayable, multi-scheduler coherence best-effort Relaxed,
off-CPU blame punted, Instant::now() uncorrected.

Zero-cost with the feature off; clippy -D warnings clean both ways;
full suite green with and without smarm-causal.
2026-07-12 19:10:04 +00:00

178 lines
6.7 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Causal-profiling demo (RFC 007): a pipeline where conventional profiling
//! lies and causal profiling doesn't.
//!
//! producer --(serialize ~200µs/item)--> reserve --(~400µs/item)--> notify
//! background: an actor burning CPU constantly, fully off the critical path
//!
//! `reserve` is the true bottleneck. `serialize` is hot but overlapped with
//! `reserve`'s backlog, and `background` is the hottest code in the process
//! while contributing nothing to throughput. A cycle profiler ranks them
//! background > reserve ≈ 2×serialize; the causal report instead shows
//! throughput responding to virtual speedups of `reserve` and (near-)ignoring
//! `serialize` and `background`.
//!
//! Stage cost is fixed *work* (a calibrated arithmetic loop), not fixed wall
//! time. This matters: a timed busy-wait absorbs injected causal delay into
//! its own budget and finishes on schedule regardless, making every
//! experiment read as a no-op (found live on a 24-core run: dead-flat
//! deltas). Real workloads are work-shaped, so the demo must be too.
//!
//! Run:
//! cargo run --release --example causal_pipeline --features smarm-causal
//!
//! Prints a summary, writes `profile.coz` (Coz plot-compatible), and — given
//! enough cores for the pipeline to actually run in parallel — checks the
//! expected separation and exits nonzero if it doesn't hold, so a CI box can
//! run this as a smoke test.
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
/// LCG-mix `iters` times in dependent sequence (unvectorizable, un-elidable),
/// staying preemptible — and causal-sampleable/delayable — via `check!()`.
fn work_iters(iters: u64) {
let mut acc = 0x2545_f491_4f6c_dd1du64;
let mut i = 0u64;
while i < iters {
let chunk_end = (i + 256).min(iters);
while i < chunk_end {
acc = acc.wrapping_mul(6364136223846793005).wrapping_add(i);
i += 1;
}
std::hint::black_box(acc);
smarm::check!();
}
}
/// Measure how many `work_iters` iterations fit in a microsecond on this
/// machine, so stage costs below are meaningful in time while staying
/// work-shaped.
fn calibrate_iters_per_us() -> u64 {
let n = 8_000_000u64;
let t = Instant::now();
work_iters(n);
(n / (t.elapsed().as_micros().max(1) as u64)).max(1)
}
fn main() {
let per_us = calibrate_iters_per_us();
println!("calibration: {per_us} work iters/µs");
let work_us = move |us: u64| work_iters(us * per_us);
let mut failures: Vec<String> = Vec::new();
smarm::init(smarm::Config::default()).run(move || {
let stop = Arc::new(AtomicBool::new(false));
let (tx_ab, rx_ab) = smarm::channel::<u64>();
let (tx_bc, rx_bc) = smarm::channel::<u64>();
// Producer: hot serialization, but upstream of the bottleneck.
let stop_p = stop.clone();
let producer = smarm::spawn(move || {
let mut i = 0u64;
while !stop_p.load(Ordering::Relaxed) {
{
let _g = smarm::causal_site!("serialize");
work_us(200);
}
if tx_ab.send(i).is_err() {
break;
}
i += 1;
}
// tx_ab drops here; downstream drains and exits.
});
// Reserve: the true bottleneck (~400µs of work per item).
let reserve = smarm::spawn(move || {
while let Ok(item) = rx_ab.recv() {
{
let _g = smarm::causal_site!("reserve");
work_us(400);
}
if tx_bc.send(item).is_err() {
break;
}
}
});
// Notify: light tail stage; marks the unit of useful work.
let notify = smarm::spawn(move || {
while rx_bc.recv().is_ok() {
{
let _g = smarm::causal_site!("notify");
work_us(50);
}
smarm::progress!("orders-processed");
}
});
// Background: hottest code in the process, zero throughput relevance.
let stop_bg = stop.clone();
let background = smarm::spawn(move || {
while !stop_bg.load(Ordering::Relaxed) {
let _g = smarm::causal_site!("background-compaction");
work_us(500);
}
});
// Warm up so queues reach steady state before measuring.
smarm::sleep(Duration::from_millis(300));
let results = smarm::causal::run_experiments(&smarm::causal::ExperimentPlan {
speedups_pct: vec![0, 25, 50],
experiment: Duration::from_millis(700),
cooldown: Duration::from_millis(150),
});
stop.store(true, Ordering::Relaxed);
producer.join().unwrap();
reserve.join().unwrap();
notify.join().unwrap();
background.join().unwrap();
print!("{}", smarm::causal::render_summary(&results));
let coz = smarm::causal::render_coz(&results);
match std::fs::write("profile.coz", coz) {
Ok(()) => println!("\nwrote profile.coz"),
Err(e) => eprintln!("\nfailed to write profile.coz: {e}"),
}
// Verdict. The separation only exists when the four pipeline actors
// actually run in parallel; on a small box, report and skip.
let cores = std::thread::available_parallelism().map(|n| n.get()).unwrap_or(1);
if cores < 4 {
println!("verdict: SKIPPED ({cores} cores; separation needs the stages in parallel)");
return;
}
let impact = |site: &str| {
smarm::causal::impact_pct(&results, site, 25, "orders-processed")
};
let mut expect = |site: &str, ok: &dyn Fn(f64) -> bool, want: &str| match impact(site) {
Some(p) => {
let verdict = if ok(p) { "ok" } else { "FAIL" };
println!("verdict: {site} @25% -> {p:+.1}% (want {want}) {verdict}");
if !ok(p) {
failures.push(format!("{site}: {p:+.1}% (want {want})"));
}
}
None => {
println!("verdict: {site} @25% -> missing cell FAIL");
failures.push(format!("{site}: missing cell"));
}
};
expect("reserve", &|p| p > 15.0, "> +15%");
expect("serialize", &|p| p < 10.0, "< +10%");
expect("background-compaction", &|p| p < 10.0, "< +10%");
if failures.is_empty() {
println!("verdict: PASS — causal separation holds");
} else {
println!("verdict: FAIL — {}", failures.join("; "));
std::process::exit(1);
}
});
}