cranpose-core 0.9.13

Composition runtime and reactive state for Cranpose
Documentation
use std::{
    future::Future,
    pin::Pin,
    sync::{
        Arc, Condvar, Mutex,
        atomic::{AtomicBool, AtomicUsize, Ordering},
    },
    task::{Context, Poll, Wake, Waker},
    time::{Duration, Instant},
};

use cranpose_core::{Composition, MemoryApplier};

#[path = "support/blocking.rs"]
mod support;
use support::Failure;

type Gate = Arc<(Mutex<bool>, Condvar)>;

fn wait(gate: &Gate) {
    let (lock, wake) = &**gate;
    let guard = lock.lock().expect("gate");
    drop(
        wake.wait_while(guard, |released| !*released)
            .expect("gate wait"),
    );
}

fn release(gate: &Gate) {
    *gate.0.lock().expect("gate") = true;
    gate.1.notify_all();
}

struct Ready {
    runnable: AtomicBool,
    waiter: std::thread::Thread,
}
impl Wake for Ready {
    fn wake(self: Arc<Self>) {
        self.runnable.store(true, Ordering::Release);
        self.waiter.unpark();
    }
    fn wake_by_ref(self: &Arc<Self>) {
        self.runnable.store(true, Ordering::Release);
        self.waiter.unpark();
    }
}

struct Slot {
    future: Pin<Box<dyn Future<Output = Result<u64, Failure>>>>,
    ready: Arc<Ready>,
}
impl Slot {
    fn new(work: impl FnOnce() -> u64 + Send + 'static) -> Self {
        Self {
            future: Box::pin(support::run(work)),
            ready: Arc::new(Ready {
                runnable: AtomicBool::new(true),
                waiter: std::thread::current(),
            }),
        }
    }
    fn poll(&mut self) -> Option<Result<u64, Failure>> {
        if !self.ready.runnable.swap(false, Ordering::AcqRel) {
            return None;
        }
        let waker = Waker::from(Arc::clone(&self.ready));
        match self.future.as_mut().poll(&mut Context::from_waker(&waker)) {
            Poll::Ready(value) => Some(value),
            Poll::Pending => None,
        }
    }
}

fn process() -> String {
    let text = std::fs::read_to_string("/proc/self/status").unwrap_or_default();
    let value = |name: &str| {
        text.lines()
            .find_map(|line| line.strip_prefix(name))
            .and_then(|line| line.split_whitespace().next())
            .unwrap_or("null")
    };
    format!(
        "{{\"rss_kib\":{},\"peak_rss_kib\":{},\"threads\":{}}}",
        value("VmRSS:"),
        value("VmHWM:"),
        value("Threads:")
    )
}

fn compute(mut value: u64) -> u64 {
    for _ in 0..1_000_000 {
        value = std::hint::black_box(
            value
                .wrapping_mul(6_364_136_223_846_793_005)
                .rotate_left(13)
                ^ 1_442_695_040_888_963_407,
        );
    }
    value
}

#[derive(Default)]
struct Jobs {
    slots: Vec<Slot>,
    next: usize,
    completed: usize,
    rejected: usize,
    checksum: u64,
}
impl Jobs {
    fn fill(&mut self, total: usize, gate: &Gate) {
        while self.next < total {
            let value = self.next as u64;
            let gate = Arc::clone(gate);
            let mut slot = Slot::new(move || {
                wait(&gate);
                compute(value)
            });
            match slot.poll() {
                Some(Err(Failure::Saturated)) => {
                    self.rejected += 1;
                    return;
                }
                Some(Err(Failure::Other)) => panic!("worker failed"),
                Some(Ok(value)) => self.complete(value),
                None => self.slots.push(slot),
            }
            self.next += 1;
        }
    }
    fn complete(&mut self, value: u64) {
        self.completed += 1;
        self.checksum = self.checksum.wrapping_add(value);
    }
    fn collect(&mut self) -> bool {
        let mut changed = false;
        for index in (0..self.slots.len()).rev() {
            if let Some(result) = self.slots[index].poll() {
                self.complete(result.expect("admitted work succeeds"));
                self.slots.swap_remove(index);
                changed = true;
            }
        }
        changed
    }
}

const GAP_BUCKETS_US: [u128; 12] = [
    1,
    4,
    16,
    64,
    100,
    250,
    500,
    1_000,
    4_000,
    16_000,
    64_000,
    u128::MAX,
];
struct Gaps {
    counts: [u64; 12],
    maximum_ns: u128,
}
impl Gaps {
    fn record(&mut self, elapsed: Duration) {
        let index = GAP_BUCKETS_US
            .iter()
            .position(|bound| elapsed.as_micros() <= *bound)
            .expect("last bucket");
        self.counts[index] += 1;
        self.maximum_ns = self.maximum_ns.max(elapsed.as_nanos());
    }
    fn percentile(&self) -> u128 {
        let target = self
            .counts
            .iter()
            .sum::<u64>()
            .saturating_mul(99)
            .div_ceil(100);
        let mut count = 0;
        for (index, amount) in self.counts.iter().enumerate() {
            count += amount;
            if count >= target {
                return GAP_BUCKETS_US[index];
            }
        }
        self.maximum_ns.div_ceil(1_000)
    }
}

fn cpu(busy_ui: bool) {
    let gate = Arc::new((Mutex::new(false), Condvar::new()));
    let mut jobs = Jobs::default();
    let composition = busy_ui.then(|| Composition::new(MemoryApplier::new()));
    let runtime = composition.as_ref().map(Composition::runtime_handle);
    let mut gaps = Gaps {
        counts: [0; 12],
        maximum_ns: 0,
    };
    let started = Instant::now();
    jobs.fill(96, &gate);
    let admission_ns = started.elapsed().as_nanos();
    let submitted = jobs.next;
    let loaded = process();
    release(&gate);
    let mut previous = Instant::now();
    while jobs.completed < 96 {
        let now = Instant::now();
        gaps.record(now.duration_since(previous));
        previous = now;
        if let Some(runtime) = &runtime {
            runtime.drain_ui();
        }
        if jobs.collect() {
            jobs.fill(96, &gate);
        } else if !busy_ui {
            std::thread::park_timeout(Duration::from_secs(60).saturating_sub(started.elapsed()));
        }
        assert!(
            started.elapsed() < Duration::from_secs(60),
            "CPU jobs stalled"
        );
    }
    let elapsed_ns = started.elapsed().as_nanos();
    let phase = if busy_ui { "contended" } else { "cpu" };
    let p99 = if busy_ui {
        gaps.percentile().to_string()
    } else {
        "null".into()
    };
    let maximum = if busy_ui {
        gaps.maximum_ns.to_string()
    } else {
        "null".into()
    };
    println!(
        "{{\"phase\":\"{phase}\",\"jobs\":{},\"checksum\":{},\"elapsed_ns\":{elapsed_ns},\"initial_admission_ns\":{admission_ns},\"initial_submitted\":{submitted},\"rejected_attempts\":{},\"main_p99_gap_us_upper\":{p99},\"main_max_gap_ns\":{maximum},\"loaded\":{loaded},\"finished\":{}}}",
        jobs.completed,
        jobs.checksum,
        jobs.rejected,
        process()
    );
}

struct Payload {
    bytes: Vec<u8>,
    live: Arc<AtomicUsize>,
}
impl Payload {
    fn new(live: &Arc<AtomicUsize>) -> Self {
        let bytes = vec![17; 128 * 1024];
        live.fetch_add(bytes.len(), Ordering::SeqCst);
        Self {
            bytes,
            live: Arc::clone(live),
        }
    }
}
impl Drop for Payload {
    fn drop(&mut self) {
        self.live.fetch_sub(self.bytes.len(), Ordering::SeqCst);
    }
}

fn cancel() {
    let gate = Arc::new((Mutex::new(false), Condvar::new()));
    let live = Arc::new(AtomicUsize::new(0));
    let executed = Arc::new(AtomicUsize::new(0));
    let mut slots = Vec::new();
    let mut rejected = 0;
    for _ in 0..128 {
        let payload = Payload::new(&live);
        let gate = Arc::clone(&gate);
        let executed = Arc::clone(&executed);
        let mut slot = Slot::new(move || {
            wait(&gate);
            executed.fetch_add(1, Ordering::SeqCst);
            std::hint::black_box(payload.bytes.len() as u64)
        });
        match slot.poll() {
            None => slots.push(slot),
            Some(Err(Failure::Saturated)) => rejected += 1,
            Some(_) => panic!("a gated job cannot finish"),
        }
    }
    let pending = slots.len();
    let before_cancel = live.load(Ordering::SeqCst);
    let loaded = process();
    let started = Instant::now();
    drop(slots);
    let cancellation_ns = started.elapsed().as_nanos();
    let after_cancel = live.load(Ordering::SeqCst);
    release(&gate);
    while live.load(Ordering::SeqCst) != 0 {
        assert!(
            started.elapsed() < Duration::from_secs(60),
            "cancelled payloads retained"
        );
        std::thread::yield_now();
    }
    println!(
        "{{\"phase\":\"cancel\",\"submitted\":128,\"pending_futures\":{pending},\"rejected\":{rejected},\"executed\":{},\"payload_before_cancel\":{before_cancel},\"payload_after_cancel\":{after_cancel},\"cancellation_ns\":{cancellation_ns},\"loaded\":{loaded},\"finished\":{}}}",
        executed.load(Ordering::SeqCst),
        process()
    );
}

fn main() {
    let mode = std::env::args()
        .skip(1)
        .find(|argument| !argument.starts_with('-'))
        .unwrap_or_else(|| "serial".into());
    match mode.as_str() {
        "cpu" => cpu(false),
        "contended" => cpu(true),
        "cancel" => cancel(),
        "serial" => {
            support::serial(1);
            let started = Instant::now();
            let checksum = support::serial(1_000);
            let elapsed_ns = started.elapsed().as_nanos();
            println!(
                "{{\"phase\":\"serial\",\"jobs\":1000,\"checksum\":{checksum},\"elapsed_ns\":{elapsed_ns},\"finished\":{}}}",
                process()
            );
        }
        _ => panic!("expected serial, cpu, contended or cancel"),
    }
}