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 {
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();
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()))
}