saddle-framework 0.3.27

The single business-facing facade for Saddle applications
//! Feature-gated, file-controlled physical holder for isolated capacity tests.
//! The production build has no hook, environment parsing or allocation path.

use std::{alloc::Layout, fs::OpenOptions, io::Write, mem::MaybeUninit, path::PathBuf};

use saddle_admission::{StorageDemand, StoredValue};
use saddle_runtime::profusegw::ProfuseGwProcessLease;
use serde::Deserialize;

#[derive(Deserialize)]
struct Command {
    seq: u64,
    action: String,
    leave_bytes: Option<usize>,
}

struct PhysicalBlock {
    bytes: Option<Box<[MaybeUninit<u8>]>>,
    layout: Layout,
    events: PathBuf,
    seq: u64,
}

impl PhysicalBlock {
    fn allocate(layout: Layout, events: PathBuf, seq: u64) -> Self {
        // This allocation follows the original process-domain permit. The
        // boxed uninitialized slice uses exactly len bytes with byte alignment.
        let mut bytes = Box::<[u8]>::new_uninit_slice(layout.size());
        for offset in (0..layout.size()).step_by(4096) {
            bytes[offset].write(0xa5);
        }
        Self { bytes: Some(bytes), layout, events, seq }
    }
}

impl Drop for PhysicalBlock {
    fn drop(&mut self) {
        let pointer = self.bytes.as_ref().map_or(0, |bytes| bytes.as_ptr() as usize);
        drop(self.bytes.take());
        append(&self.events, serde_json::json!({"seq":self.seq,"event":"physical_destroyed",
            "address":pointer,"layout_bytes":self.layout.size(),"layout_align":self.layout.align()}));
    }
}

pub(super) struct Controller {
    control: PathBuf,
    events: PathBuf,
    last_seq: u64,
    held: Option<Held>,
}

enum Held {
    Framework(StoredValue<PhysicalBlock>),
    Managed(saddle_admission::TestManagedHolder),
}

impl Controller {
    pub(super) fn from_env() -> Option<Self> {
        let control = PathBuf::from(std::env::var_os("SADDLE_TEST_HOLDER_CONTROL")?);
        let events = PathBuf::from(std::env::var_os("SADDLE_TEST_HOLDER_EVENTS")?);
        if !control.is_absolute() || !events.is_absolute() { return None; }
        Some(Self { control, events, last_seq: 0, held: None })
    }

    pub(super) fn poll(&mut self, lease: &ProfuseGwProcessLease, stop: &std::sync::atomic::AtomicBool) {
        let Ok(bytes) = std::fs::read(&self.control) else { return; };
        let Ok(command) = serde_json::from_slice::<Command>(&bytes) else { return; };
        if command.seq <= self.last_seq { return; }
        self.last_seq = command.seq;
        let before = lease.resource_snapshot();
        let result = match command.action.as_str() {
            "hold" => self.hold(lease, &command),
            "hold_managed" => self.hold_managed(lease, &command),
            "inspect" => {
                append(&self.events, serde_json::json!({"seq":command.seq,
                    "event":"ownership", "snapshot":lease.ownership_snapshot(),
                    "response_capacity":crate::profusegw_http::test_response_capacity_snapshot()}));
                Ok(())
            }
            "shutdown" => {
                stop.store(true, std::sync::atomic::Ordering::Release);
                Ok(())
            }
            "release" => self.release(),
            _ => Err("unknown_action"),
        };
        let after = lease.resource_snapshot();
        append(&self.events, serde_json::json!({"seq":command.seq,"event":"command",
            "action":command.action,"result":result.err(),
            "rss_kib":rss_kib(),
            "before_framework_charged":before.map(|x|x.framework_charged),
            "after_framework_charged":after.map(|x|x.framework_charged),
            "after_framework_available":after.map(|x|x.framework_capacity.saturating_sub(x.framework_charged)),
            "before_managed_committed":before.map(|x|x.committed),
            "after_managed_committed":after.map(|x|x.committed),
            "after_managed_available":after.map(|x|x.managed_capacity.saturating_sub(x.committed))}));
    }

    pub(super) fn record_supervisor_drop(
        &self,
        before: Option<saddle_admission::ProcessOwnershipSnapshot>,
        after: Option<saddle_admission::ProcessOwnershipSnapshot>,
    ) {
        append(&self.events, serde_json::json!({"event":"supervisor_drained",
            "before":before,"after":after}));
    }

    fn hold(&mut self, lease: &ProfuseGwProcessLease, command: &Command) -> Result<(), &'static str> {
        if self.held.is_some() { return Err("already_held"); }
        let leave = command.leave_bytes.ok_or("missing_leave_bytes")?;
        let snapshot = lease.resource_snapshot().ok_or("process_closed")?;
        let available = snapshot.framework_capacity.saturating_sub(snapshot.framework_charged);
        let bytes = available.checked_sub(leave).and_then(|x|x.checked_sub(std::mem::size_of::<saddle_admission::StoragePermit>()))
            .ok_or("insufficient_available")?;
        let layout = Layout::from_size_align(bytes, 1).map_err(|_|"layout_invalid")?;
        if bytes == 0 { return Err("empty_layout"); }
        let demand = StorageDemand::separate(&[(layout,1)]).map_err(|_|"demand_invalid")?;
        let demand_bytes = demand.bytes();
        let permit = lease.try_test_process_storage(demand).map_err(|_|"original_domain_rejected")?;
        let block = PhysicalBlock::allocate(layout, self.events.clone(), command.seq);
        let pointer = block.bytes.as_ref().unwrap().as_ptr() as usize;
        self.held = Some(Held::Framework(permit.hold(block)));
        append(&self.events, serde_json::json!({"seq":command.seq,"event":"physical_held",
            "address":pointer,"layout_bytes":layout.size(),"layout_align":layout.align(),
            "demand_bytes":demand_bytes,"rss_kib":rss_kib()}));
        Ok(())
    }

    fn hold_managed(&mut self, lease: &ProfuseGwProcessLease, command: &Command) -> Result<(), &'static str> {
        if self.held.is_some() { return Err("already_held"); }
        let leave = command.leave_bytes.ok_or("missing_leave_bytes")?;
        let snapshot = lease.resource_snapshot().ok_or("process_closed")?;
        let available = snapshot.managed_capacity.saturating_sub(snapshot.committed);
        let bytes = available.checked_sub(leave).ok_or("insufficient_available")?;
        if bytes == 0 { return Err("empty_layout"); }
        let mut holder = lease.try_test_managed_holder(bytes).map_err(|_|"original_managed_domain_rejected")?;
        let address = holder.address();
        // Abort this isolated test object before the frozen 85% observation
        // waterline. The original 4GiB cgroup is an additional hard backstop.
        let safety_kib = ((4_u64 * 1024 * 1024 * 1024 * 85) / 100 / 1024) as usize;
        for offset in (0..bytes).step_by(4096) {
            holder.touch_page(offset);
            if offset % (16 * 1024 * 1024) == 0 && rss_kib().is_some_and(|rss| rss >= safety_kib) {
                return Err("rss_safety_waterline");
            }
        }
        holder.touch_page(bytes - 1);
        self.held = Some(Held::Managed(holder));
        append(&self.events, serde_json::json!({"seq":command.seq,"event":"physical_held",
            "domain":"managed","address":address,"layout_bytes":bytes,"layout_align":1,
            "demand_bytes":bytes,"rss_kib":rss_kib()}));
        Ok(())
    }

    fn release(&mut self) -> Result<(), &'static str> {
        match self.held.take() {
            Some(Held::Framework(block)) => { drop(block); Ok(()) },
            Some(Held::Managed(block)) => {
                let address = block.address(); let capacity = block.capacity();
                block.release(|| append(&self.events, serde_json::json!({
                    "event":"physical_destroyed","domain":"managed","address":address,
                    "layout_bytes":capacity,"layout_align":1,"rss_kib":rss_kib()
                }))).map_err(|_|"managed_release_failed")
            },
            None => Err("not_held"),
        }
    }
}

fn append(path: &PathBuf, value: serde_json::Value) {
    if let Ok(mut file) = OpenOptions::new().create(true).append(true).open(path) {
        let _ = writeln!(file,"{value}");
        let _ = file.flush();
    }
}

fn rss_kib() -> Option<usize> {
    let status=std::fs::read_to_string("/proc/self/status").ok()?;
    status.lines().find_map(|line|line.strip_prefix("VmRSS:")
        .and_then(|value|value.split_ascii_whitespace().next())
        .and_then(|value|value.parse().ok()))
}