use super::cgfs;
use super::journal::{should_restore, Journal, JournalAction, JournalEntry};
use super::resolve::{Mechanism, Resolution};
use super::sampler::parse_meminfo;
use super::systemd::SystemdUser;
use super::types::Action;
use crate::CgroupManager;
use common::Result;
use std::process::Command;
use std::time::Duration;
pub const MIN_CAP_BYTES: u64 = 256 * 1024 * 1024;
const DBUS_TIMEOUT: Duration = Duration::from_secs(2);
pub struct Effector<'a> {
manager: &'a CgroupManager,
journal: &'a Journal,
systemd: Option<&'a SystemdUser>,
}
impl<'a> Effector<'a> {
pub fn new(
manager: &'a CgroupManager,
journal: &'a Journal,
systemd: Option<&'a SystemdUser>,
) -> Self {
Self {
manager,
journal,
systemd,
}
}
pub fn apply(&self, action: &Action) -> Result<()> {
match action {
Action::Freeze { res, name } => self.freeze(res, name),
Action::Thaw { res } => self.thaw(res),
Action::Cap { res, name } => self.cap(res, name),
Action::LiftCap { res } => self.lift_cap(res),
Action::Notify { message } => {
notify(message);
Ok(())
}
}
}
fn freeze(&self, res: &Resolution, name: &str) -> Result<()> {
let Some(inode) = cgfs::dir_inode(&res.cgroup) else {
tracing::warn!(
cgroup = %res.cgroup, name,
"cannot read cgroup inode; refusing to freeze (would be unrestorable)"
);
return Err(common::Error::Cgroup(format!(
"cannot read inode for {}; refusing to freeze",
res.cgroup
)));
};
let entry = JournalEntry {
cgroup: res.cgroup.clone(),
inode,
unit: res.unit.clone(),
action: JournalAction::Freeze,
prev_high: None,
our_high: None,
};
self.journal.append(&entry)?;
tracing::info!(cgroup = %res.cgroup, name, "freezing cgroup");
if res.mechanism == Mechanism::Unit {
if let (Some(unit), Some(systemd)) = (&res.unit, self.systemd) {
match systemd.freeze_unit(unit, DBUS_TIMEOUT) {
Ok(()) => return Ok(()),
Err(e) => tracing::warn!(
cgroup = %res.cgroup, unit, error = %e,
"FreezeUnit failed; falling back to raw cgroup.freeze"
),
}
}
}
cgfs::write_freeze(&res.cgroup, true)
}
fn thaw(&self, res: &Resolution) -> Result<()> {
tracing::info!(cgroup = %res.cgroup, "thawing cgroup");
let entries = self.entries_for(&res.cgroup);
let result = self.thaw_raw(&res.cgroup, res.unit.as_deref());
if let Err(e) = &result {
tracing::debug!(cgroup = %res.cgroup, error = %e, "raw thaw failed (cgroup may already be gone)");
}
self.restore_high_if_any(&res.cgroup, &entries);
self.journal.remove(&res.cgroup)?;
result
}
fn cap(&self, res: &Resolution, name: &str) -> Result<()> {
let Some(inode) = cgfs::dir_inode(&res.cgroup) else {
tracing::warn!(
cgroup = %res.cgroup, name,
"cannot read cgroup inode; refusing to cap (would be unrestorable)"
);
return Err(common::Error::Cgroup(format!(
"cannot read inode for {}; refusing to cap",
res.cgroup
)));
};
let prev_high = cgfs::read_high(&res.cgroup);
let swap_total_kb = std::fs::read_to_string("/proc/meminfo")
.ok()
.and_then(|m| parse_meminfo(&m))
.map_or(0, |m| m.swap_total_kb);
let anon_ok = cgfs::anon_reclaimable(&res.cgroup, swap_total_kb);
let Some(target) = cap_target(
cgfs::current_bytes(&res.cgroup),
cgfs::file_bytes(&res.cgroup),
anon_ok,
) else {
tracing::warn!(
cgroup = %res.cgroup, name,
"cannot read memory.current; refusing to cap"
);
return Err(common::Error::Cgroup(
"cannot read memory.current; refusing to cap".into(),
));
};
let our_bytes = page_align_down(target, page_size());
if !cap_tightens(prev_high.as_deref(), our_bytes) {
tracing::info!(
cgroup = %res.cgroup, name, prev_high = ?prev_high, our_bytes,
"existing memory.high already at or below the cap; not capping"
);
return Err(common::Error::Cgroup(
"existing memory.high already at or below the cap; refusing to cap".into(),
));
}
let our_high = our_bytes.to_string();
let entry = JournalEntry {
cgroup: res.cgroup.clone(),
inode,
unit: res.unit.clone(),
action: JournalAction::Cap,
prev_high,
our_high: Some(our_high.clone()),
};
self.journal.append(&entry)?;
tracing::info!(
cgroup = %res.cgroup, name, our_high = %our_high, anon_reclaimable = anon_ok,
"soft-capping cgroup"
);
let result = cgfs::write_high(&res.cgroup, &our_high);
if result.is_ok() {
self.reconcile_our_high(&entry);
}
result
}
fn reconcile_our_high(&self, written: &JournalEntry) {
let cgroup: &str = &written.cgroup;
let Some(actual) = cgfs::read_high(cgroup) else {
return;
};
if written.our_high.as_deref() == Some(actual.as_str()) {
return;
}
tracing::warn!(
cgroup = %cgroup, journaled = ?written.our_high, actual = %actual,
"memory.high on disk differs from what we journaled; correcting journal entry"
);
let mut entries = self.entries_for(cgroup);
let Some(pos) = entries.iter().rposition(|e| e == written) else {
return;
};
entries[pos].our_high = Some(actual);
if let Err(e) = self.journal.replace(cgroup, &entries) {
tracing::warn!(cgroup = %cgroup, error = %e, "failed to correct journal entry (atomic replace)");
}
}
fn lift_cap(&self, res: &Resolution) -> Result<()> {
tracing::info!(cgroup = %res.cgroup, "lifting cap");
let entries = self.entries_for(&res.cgroup);
let _ = self.thaw_raw(&res.cgroup, res.unit.as_deref());
self.restore_high_if_any(&res.cgroup, &entries);
self.journal.remove(&res.cgroup)
}
fn entries_for(&self, cgroup: &str) -> Vec<JournalEntry> {
self.journal
.entries()
.into_iter()
.filter(|e| e.cgroup == cgroup)
.collect()
}
pub fn sweep_leftovers(&self) -> Result<()> {
if let Err(e) = self.manager.sweep_guard_leftovers() {
tracing::warn!(error = %e, "legacy guard-<pid> sweep failed (non-fatal)");
}
self.replay_and_clear()
}
pub fn undo_all(&self) -> Result<()> {
self.replay_and_clear()
}
fn replay_and_clear(&self) -> Result<()> {
let mut by_cgroup: std::collections::HashMap<String, Vec<JournalEntry>> =
std::collections::HashMap::new();
for e in self.journal.entries() {
by_cgroup.entry(e.cgroup.clone()).or_default().push(e);
}
for (cgroup, entries) in &by_cgroup {
if let Err(e) = cgfs::write_freeze(cgroup, false) {
tracing::debug!(cgroup, error = %e, "raw thaw failed (cgroup may already be gone)");
}
self.restore_high_if_any(cgroup, entries);
}
let cleared = self.journal.clear();
if let Some(systemd) = self.systemd {
for (cgroup, entries) in &by_cgroup {
if let Some(unit) = entries.last().and_then(|e| e.unit.as_deref()) {
if let Err(e) = systemd.thaw_unit(unit, DBUS_TIMEOUT) {
tracing::debug!(cgroup, unit, error = %e, "ThawUnit failed after raw thaw");
}
}
}
}
cleared
}
fn restore_high_if_any(&self, cgroup: &str, entries: &[JournalEntry]) {
let inode = cgfs::dir_inode(cgroup);
let high = cgfs::read_high(cgroup);
match restore_decision(entries, inode, high.as_deref()) {
RestoreStep::ThawAndRestoreHigh { to } => {
if let Err(err) = cgfs::write_high(cgroup, &to) {
tracing::warn!(cgroup, error = %err, "failed to restore memory.high");
}
}
RestoreStep::ThawOnly => {}
RestoreStep::SkipRemove => {
tracing::warn!(
cgroup,
"not restoring memory.high (cgroup recreated or value changed since our write)"
);
}
}
}
fn thaw_raw(&self, cgroup: &str, unit: Option<&str>) -> Result<()> {
let result = cgfs::write_freeze(cgroup, false);
if let (Some(unit), Some(systemd)) = (unit, self.systemd) {
if let Err(e) = systemd.thaw_unit(unit, DBUS_TIMEOUT) {
tracing::debug!(cgroup, unit, error = %e, "ThawUnit failed after raw thaw");
}
}
result
}
}
#[derive(Debug, PartialEq, Eq)]
pub enum RestoreStep {
ThawOnly,
ThawAndRestoreHigh { to: String },
SkipRemove,
}
pub fn restore_step(e: &JournalEntry, inode: Option<u64>, high: Option<&str>) -> RestoreStep {
if !should_restore(e, inode, high) {
return RestoreStep::SkipRemove;
}
match e.action {
JournalAction::Freeze => RestoreStep::ThawOnly,
JournalAction::Cap => RestoreStep::ThawAndRestoreHigh {
to: e.prev_high.clone().unwrap_or_else(|| "max".into()),
},
}
}
pub fn cap_target(current: Option<u64>, file: Option<u64>, anon_reclaimable: bool) -> Option<u64> {
let current = current?;
let mut target = current / 10 * 9;
if !anon_reclaimable {
let file_floor = current.saturating_sub(file.unwrap_or(0).saturating_mul(8) / 10);
target = target.max(file_floor);
}
Some(target.max(MIN_CAP_BYTES))
}
pub fn cap_tightens(prev_high: Option<&str>, target: u64) -> bool {
match prev_high.and_then(|s| s.trim().parse::<u64>().ok()) {
Some(prev) => prev > target,
None => true,
}
}
pub fn restore_decision(
entries: &[JournalEntry],
inode: Option<u64>,
high: Option<&str>,
) -> RestoreStep {
let Some(newest_cap) = entries
.iter()
.rev()
.find(|e| e.action == JournalAction::Cap)
else {
return RestoreStep::ThawOnly;
};
match restore_step(newest_cap, inode, high) {
RestoreStep::ThawAndRestoreHigh { .. } => match restore_target(entries) {
Some(to) => RestoreStep::ThawAndRestoreHigh { to },
None => RestoreStep::ThawOnly,
},
other => other,
}
}
pub fn restore_target(entries: &[JournalEntry]) -> Option<String> {
entries
.iter()
.find(|e| e.action == JournalAction::Cap)
.map(|e| e.prev_high.clone().unwrap_or_else(|| "max".into()))
}
fn page_size() -> u64 {
let p = unsafe { libc::sysconf(libc::_SC_PAGESIZE) };
if p > 0 {
p as u64
} else {
4096
}
}
fn page_align_down(bytes: u64, page: u64) -> u64 {
bytes.checked_div(page).map_or(bytes, |q| q * page)
}
fn notify(message: &str) {
match Command::new("notify-send")
.arg("rlm-guard")
.arg(message)
.spawn()
{
Ok(mut child) => {
std::thread::spawn(move || {
let _ = child.wait();
});
}
Err(e) => {
tracing::debug!(error = %e, "notify-send unavailable; skipping notification");
}
}
}
#[cfg(test)]
mod tests {
use super::super::resolve::{Coverage, Verdict};
use super::*;
const MIB: u64 = 1024 * 1024;
const GIB: u64 = 1024 * MIB;
#[test]
fn cap_never_demands_more_than_ten_percent_of_current() {
let current = GIB * 3 / 2;
let cap = cap_target(Some(current), Some(GIB * 135 / 100), true).unwrap();
assert!(
cap >= current / 10 * 9,
"cap {cap} is below 90% of {current}"
);
}
#[test]
fn swapless_cap_only_asks_for_reclaimable_file_pages() {
assert_eq!(
cap_target(Some(GIB), Some(0), false),
Some(GIB),
"nothing reclaimable: demand nothing"
);
assert_eq!(
cap_target(Some(GIB), Some(512 * MIB), false),
Some(GIB / 10 * 9)
);
assert_eq!(
cap_target(Some(GIB), None, false),
Some(GIB),
"unknown file size: demand nothing"
);
}
#[test]
fn cap_has_a_256_mib_floor() {
assert_eq!(
cap_target(Some(100 * MIB), Some(0), true),
Some(MIN_CAP_BYTES)
);
}
#[test]
fn cap_never_loosens_an_existing_memory_high() {
let floor = cap_target(Some(100 * MIB), Some(0), true).unwrap();
assert_eq!(floor, MIN_CAP_BYTES);
let prev = (200 * MIB).to_string();
assert!(
!cap_tightens(Some(&prev), floor),
"200 MiB unit limit must not be raised to the 256 MiB floor"
);
assert!(
!cap_tightens(Some(&floor.to_string()), floor),
"equal: no-op"
);
assert!(cap_tightens(Some("max"), floor));
assert!(cap_tightens(None, floor));
let prev = (2 * GIB).to_string();
assert!(cap_tightens(Some(&prev), GIB));
}
#[test]
fn unreadable_current_refuses_to_cap() {
assert_eq!(cap_target(None, Some(1), true), None);
}
#[test]
fn our_high_string_is_plain_decimal_no_separators() {
let bytes = cap_target(Some(12_345_678_900), None, true).unwrap();
let s = bytes.to_string();
assert!(
s.chars().all(|c| c.is_ascii_digit()),
"our_high must be plain digits, got {s:?}"
);
assert_eq!(s, format!("{bytes}"), "no formatting beyond plain decimal");
}
#[test]
fn page_align_down_rounds_to_page_multiple() {
assert_eq!(page_align_down(900_000_000, 4096), 899_997_696);
assert_eq!(
page_align_down(4096, 4096),
4096,
"already-aligned is a no-op"
);
assert_eq!(page_align_down(100, 0), 100, "page=0 guard is a no-op");
}
#[test]
fn restore_target_uses_oldest_caps_prev_high() {
let oldest = entry_cap("/x", 42, "max", "A");
let newest = entry_cap("/x", 42, "A-prime", "B");
assert_eq!(
restore_target(&[oldest, newest]),
Some("max".into()),
"must use the oldest entry's prev_high, not the newest's"
);
}
#[test]
fn restore_target_none_for_freeze_only_chain() {
assert_eq!(restore_target(&[entry_freeze("/x", 42)]), None);
}
#[test]
fn restore_decision_liveness_uses_newest_cap_value_uses_oldest() {
let oldest = entry_cap("/x", 42, "max", "A");
let newest = entry_cap("/x", 42, "A", "B");
assert_eq!(
restore_decision(&[oldest.clone(), newest.clone()], Some(42), Some("B")),
RestoreStep::ThawAndRestoreHigh { to: "max".into() },
"must judge liveness against the newest Cap's our_high (matches disk \"B\"), \
but restore the oldest Cap's prev_high (\"max\")"
);
assert_eq!(
restore_step(&oldest, Some(42), Some("B")),
RestoreStep::SkipRemove,
"confirms the bug this guards against: judging liveness against the oldest \
entry's our_high against the newer on-disk value mismatches"
);
let cap = entry_cap("/x", 42, "max", "1000");
let frz = entry_freeze("/x", 42);
assert_eq!(
restore_decision(&[cap, frz], Some(42), Some("1000")),
RestoreStep::ThawAndRestoreHigh { to: "max".into() }
);
let cap_only = entry_cap("/x", 42, "max", "1000");
assert_eq!(
restore_decision(&[cap_only], Some(42), Some("1000")),
RestoreStep::ThawAndRestoreHigh { to: "max".into() }
);
assert_eq!(
restore_decision(&[entry_freeze("/x", 42)], Some(42), None),
RestoreStep::ThawOnly
);
}
fn entry_cap(cg: &str, inode: u64, prev: &str, our: &str) -> JournalEntry {
JournalEntry {
cgroup: cg.into(),
inode,
unit: None,
action: JournalAction::Cap,
prev_high: Some(prev.into()),
our_high: Some(our.into()),
}
}
fn entry_freeze(cg: &str, inode: u64) -> JournalEntry {
JournalEntry {
cgroup: cg.into(),
inode,
unit: None,
action: JournalAction::Freeze,
prev_high: None,
our_high: None,
}
}
#[test]
fn restore_step_matrix() {
let cap = entry_cap("/x", 42, "max", "1000");
assert_eq!(
restore_step(&cap, Some(42), Some("1000")),
RestoreStep::ThawAndRestoreHigh { to: "max".into() }
);
assert_eq!(
restore_step(&cap, Some(43), Some("1000")),
RestoreStep::SkipRemove
);
assert_eq!(
restore_step(&cap, Some(42), Some("777")),
RestoreStep::SkipRemove
);
let frz = entry_freeze("/x", 42);
assert_eq!(
restore_step(&frz, Some(42), None),
RestoreStep::ThawOnly,
"alive freeze entry: nothing to restore beyond the unconditional thaw"
);
assert_eq!(
restore_step(&frz, None, None),
RestoreStep::SkipRemove,
"dead cgroup: restore_step only decides memory.high, never freeze — the \
unconditional thaw in Effector::thaw_raw already ran before this is consulted, \
so a still-frozen dead cgroup is never left behind (carry-forward: Task 5 review)"
);
}
fn wait_for_frozen(cgroup: &str, want: bool, timeout: Duration) -> Option<bool> {
let deadline = std::time::Instant::now() + timeout;
loop {
let got = cgfs::read_frozen(cgroup);
if got == Some(want) || std::time::Instant::now() >= deadline {
return got;
}
std::thread::sleep(Duration::from_millis(20));
}
}
fn test_resolution(cgroup: String) -> Resolution {
Resolution {
cgroup,
unit: None,
verdict: Verdict::Freeze,
coverage: Coverage::Full,
mechanism: Mechanism::Raw,
}
}
#[test]
#[ignore = "requires cgroup v2 delegation; run manually"]
fn freeze_thaw_real_process_raw_cgroup() {
use common::Limit;
use std::process::Command;
let manager = CgroupManager::new().expect("create CgroupManager");
let journal_dir = tempfile::tempdir().unwrap();
let journal =
Journal::open(journal_dir.path().join("j.jsonl"), "test-boot".into()).unwrap();
let effector = Effector::new(&manager, &journal, None);
let abs_path = manager
.prepare_cgroup("test-freeze-thaw", &Limit::default())
.expect("create test cgroup")
.path;
let cgroup = format!(
"/{}",
abs_path
.strip_prefix("/sys/fs/cgroup")
.expect("cgroup under /sys/fs/cgroup")
.display()
);
let mut child = Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn sleep");
let pid = child.id();
manager
.add_to_cgroup(&abs_path, pid)
.expect("add sleep to test cgroup");
let res = test_resolution(cgroup.clone());
effector
.apply(&Action::Freeze {
res: res.clone(),
name: "sleep".into(),
})
.expect("freeze");
assert_eq!(
wait_for_frozen(&cgroup, true, Duration::from_secs(2)),
Some(true),
"cgroup should be frozen"
);
assert_eq!(journal.entries().len(), 1, "freeze should be journaled");
effector
.apply(&Action::Thaw { res: res.clone() })
.expect("thaw");
assert_eq!(
wait_for_frozen(&cgroup, false, Duration::from_secs(2)),
Some(false),
"cgroup should be thawed"
);
assert!(
journal.entries().is_empty(),
"journal entry should be removed after thaw"
);
let _ = child.kill();
let _ = child.wait();
let _ = manager.cleanup_cgroup("test-freeze-thaw");
}
#[test]
#[ignore = "requires a session bus and cgroup v2 delegation; run manually"]
fn freeze_thaw_real_transient_scope() {
use std::process::Command;
let uid_out = Command::new("id").arg("-u").output().expect("id -u");
let uid: u32 = String::from_utf8_lossy(&uid_out.stdout)
.trim()
.parse()
.expect("parse uid");
let unit_base = format!("rlm-e2e-{}", std::process::id());
let unit = format!("{unit_base}.scope");
let cgroup = format!("/user.slice/user-{uid}.slice/user@{uid}.service/app.slice/{unit}");
let mut child = Command::new("systemd-run")
.args([
"--user",
"--scope",
"--slice=app.slice",
&format!("--unit={unit_base}"),
"--",
"sleep",
"30",
])
.spawn()
.expect("spawn systemd-run --scope");
std::thread::sleep(Duration::from_millis(300));
let manager = CgroupManager::new().expect("create CgroupManager");
let journal_dir = tempfile::tempdir().unwrap();
let journal =
Journal::open(journal_dir.path().join("j.jsonl"), "test-boot".into()).unwrap();
let systemd = SystemdUser::connect();
let effector = Effector::new(&manager, &journal, systemd.as_ref());
let res = Resolution {
cgroup: cgroup.clone(),
unit: Some(unit),
verdict: Verdict::Freeze,
coverage: Coverage::Full,
mechanism: Mechanism::Unit,
};
effector
.apply(&Action::Freeze {
res: res.clone(),
name: "sleep".into(),
})
.expect("freeze");
assert_eq!(
wait_for_frozen(&cgroup, true, Duration::from_secs(2)),
Some(true),
"scope cgroup should be frozen"
);
assert_eq!(journal.entries().len(), 1, "freeze should be journaled");
effector.apply(&Action::Thaw { res }).expect("thaw");
assert_eq!(
wait_for_frozen(&cgroup, false, Duration::from_secs(2)),
Some(false),
"scope cgroup should be thawed"
);
assert!(
journal.entries().is_empty(),
"journal entry should be removed after thaw"
);
let _ = child.kill();
let _ = child.wait();
let _ = Command::new("systemctl")
.args(["--user", "stop", &format!("{unit_base}.scope")])
.status();
}
#[test]
#[ignore = "requires cgroup v2 delegation; run manually"]
fn cap_page_aligns_and_journal_matches_on_disk_value() {
use common::Limit;
use std::process::Command;
let manager = CgroupManager::new().expect("create CgroupManager");
let journal_dir = tempfile::tempdir().unwrap();
let journal =
Journal::open(journal_dir.path().join("j.jsonl"), "test-boot".into()).unwrap();
let effector = Effector::new(&manager, &journal, None);
let abs_path = manager
.prepare_cgroup("test-cap-align", &Limit::default())
.expect("create test cgroup")
.path;
let cgroup = format!(
"/{}",
abs_path
.strip_prefix("/sys/fs/cgroup")
.expect("cgroup under /sys/fs/cgroup")
.display()
);
let mut child = Command::new("bash")
.arg("-c")
.arg("a=$(head -c 300000000 /dev/urandom | base64 -w0); sleep 30")
.spawn()
.expect("spawn memory-holding process");
let pid = child.id();
manager
.add_to_cgroup(&abs_path, pid)
.expect("add process to test cgroup");
let want = 300 * 1024 * 1024;
let deadline = std::time::Instant::now() + Duration::from_secs(10);
loop {
let cur = cgfs::current_bytes(&cgroup).unwrap_or(0);
if cur >= want {
break;
}
if std::time::Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
let _ = manager.cleanup_cgroup("test-cap-align");
panic!("memory.current reached only {cur} bytes, need {want}, after 10s");
}
std::thread::sleep(Duration::from_millis(100));
}
let res = test_resolution(cgroup.clone());
effector
.apply(&Action::Cap {
res: res.clone(),
name: "mem-hog".into(),
})
.expect("cap");
let entries = journal.entries();
assert_eq!(entries.len(), 1, "cap should be journaled");
let journaled = entries[0].our_high.clone().expect("our_high recorded");
let on_disk = cgfs::read_high(&cgroup).expect("read memory.high");
assert_eq!(
on_disk, journaled,
"on-disk memory.high must match the journaled our_high exactly"
);
let bytes: u64 = journaled.parse().expect("our_high is a plain decimal");
assert_eq!(bytes % page_size(), 0, "written value must be page-aligned");
effector.apply(&Action::LiftCap { res }).expect("lift cap");
assert!(
journal.entries().is_empty(),
"journal entry removed after lift"
);
let _ = child.kill();
let _ = child.wait();
let _ = manager.cleanup_cgroup("test-cap-align");
}
#[test]
#[ignore = "requires a session bus and cgroup v2 delegation; run manually"]
fn lift_cap_leaves_unit_memory_high_property_untouched() {
use std::process::Command;
let unit_memory_high = |unit: &str| -> String {
let out = Command::new("systemctl")
.args(["--user", "show", "-p", "MemoryHigh", "--value", unit])
.output()
.expect("systemctl --user show");
String::from_utf8_lossy(&out.stdout).trim().to_string()
};
let uid_out = Command::new("id").arg("-u").output().expect("id -u");
let uid: u32 = String::from_utf8_lossy(&uid_out.stdout)
.trim()
.parse()
.expect("parse uid");
let unit_base = format!("rlm-e2e-cap-{}", std::process::id());
let unit = format!("{unit_base}.scope");
let cgroup = format!("/user.slice/user-{uid}.slice/user@{uid}.service/app.slice/{unit}");
let mut child = Command::new("systemd-run")
.args([
"--user",
"--scope",
"--slice=app.slice",
&format!("--unit={unit_base}"),
"--property=MemoryHigh=1500M",
"--",
"sleep",
"30",
])
.spawn()
.expect("spawn systemd-run --scope with MemoryHigh set");
std::thread::sleep(Duration::from_millis(300));
let manager = CgroupManager::new().expect("create CgroupManager");
let journal_dir = tempfile::tempdir().unwrap();
let journal =
Journal::open(journal_dir.path().join("j.jsonl"), "test-boot".into()).unwrap();
let systemd = SystemdUser::connect();
let effector = Effector::new(&manager, &journal, systemd.as_ref());
let original_high =
cgfs::read_high(&cgroup).expect("systemd-run set an initial memory.high");
assert_ne!(
original_high, "max",
"test needs a concrete prior MemoryHigh to distinguish from systemd's clear-to-max"
);
let property_before = unit_memory_high(&unit);
let res = Resolution {
cgroup: cgroup.clone(),
unit: Some(unit.clone()),
verdict: Verdict::CapOnly,
coverage: Coverage::Full,
mechanism: Mechanism::Unit,
};
effector
.apply(&Action::Cap {
res: res.clone(),
name: "sleep".into(),
})
.expect("cap");
assert_ne!(
cgfs::read_high(&cgroup),
Some(original_high.clone()),
"cap should have changed memory.high"
);
effector.apply(&Action::LiftCap { res }).expect("lift cap");
assert_eq!(
unit_memory_high(&unit),
property_before,
"cap/lift must not change the unit's MemoryHigh property"
);
assert_eq!(
cgfs::read_high(&cgroup),
Some(original_high),
"lift must restore the raw pre-cap memory.high"
);
assert!(
journal.entries().is_empty(),
"journal entry removed after lift"
);
let _ = child.kill();
let _ = child.wait();
let _ = Command::new("systemctl")
.args(["--user", "stop", &unit])
.status();
}
#[test]
#[ignore = "requires cgroup v2 delegation; run manually"]
fn chain_restores_cap_value_when_newest_entry_is_freeze() {
use common::Limit;
use std::process::Command;
let manager = CgroupManager::new().expect("create CgroupManager");
let journal_dir = tempfile::tempdir().unwrap();
let journal =
Journal::open(journal_dir.path().join("j.jsonl"), "test-boot".into()).unwrap();
let effector = Effector::new(&manager, &journal, None);
let abs_path = manager
.prepare_cgroup("test-chain-restore", &Limit::default())
.expect("create test cgroup")
.path;
let cgroup = format!(
"/{}",
abs_path
.strip_prefix("/sys/fs/cgroup")
.expect("cgroup under /sys/fs/cgroup")
.display()
);
let mut child = Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn sleep");
let pid = child.id();
manager
.add_to_cgroup(&abs_path, pid)
.expect("add sleep to test cgroup");
let original_high = cgfs::read_high(&cgroup).expect("read initial memory.high");
let res = test_resolution(cgroup.clone());
effector
.apply(&Action::Cap {
res: res.clone(),
name: "sleep".into(),
})
.expect("cap");
assert_ne!(
cgfs::read_high(&cgroup),
Some(original_high.clone()),
"cap should have changed memory.high"
);
effector
.apply(&Action::Freeze {
res: res.clone(),
name: "sleep".into(),
})
.expect("freeze");
let entries = journal.entries();
assert_eq!(entries.len(), 2, "both Cap and Freeze entries coexist");
assert_eq!(entries[0].action, JournalAction::Cap, "Cap is oldest");
assert_eq!(entries[1].action, JournalAction::Freeze, "Freeze is newest");
effector.apply(&Action::Thaw { res }).expect("thaw");
assert_eq!(
cgfs::read_high(&cgroup),
Some(original_high),
"thaw must restore the chain's Cap value even though Freeze is the newest entry"
);
assert!(
journal.entries().is_empty(),
"both chain entries removed after thaw"
);
let _ = child.kill();
let _ = child.wait();
let _ = manager.cleanup_cgroup("test-chain-restore");
}
}