From 4150ee09600ee0f4a1c484a91aca1c2c89aa35f3 Mon Sep 17 00:00:00 2001 From: Sandro Wenzel Date: Mon, 5 Oct 2026 16:50:16 +0200 Subject: [PATCH] Count backfill memory in the memory gate of the new runner This fixes a problem in the memory admission of the new workflow runner and adds unit tests. - The default tier ignored memory booked by backfill tasks, so `--mem-limit` could be exceeded. - `fits_default` and `mem_free_default` now include the backfill memory. - The backfill memory factor now defaults to 1.0 and can be set with `--backfill-mem-factor`. - The simulator default is aligned with the runner. Co-Authored-By: Claude Sonnet 5.5 --- MC/workflow_runner/o2dpg_runner/cli.py | 3 +++ MC/workflow_runner/o2dpg_runner/config.py | 1 + MC/workflow_runner/o2dpg_runner/executor.py | 1 + MC/workflow_runner/o2dpg_runner/resources.py | 6 ++--- .../o2dpg_runner/tests/test_resources.py | 27 +++++++++++++++++++ .../o2dpg_schedule_simulator.py | 8 +++--- 6 files changed, 39 insertions(+), 7 deletions(-) diff --git a/MC/workflow_runner/o2dpg_runner/cli.py b/MC/workflow_runner/o2dpg_runner/cli.py index 222780782..79e0cfe46 100644 --- a/MC/workflow_runner/o2dpg_runner/cli.py +++ b/MC/workflow_runner/o2dpg_runner/cli.py @@ -60,6 +60,8 @@ def build_parser() -> argparse.ArgumentParser: p.add_argument("--dynamic-resources", dest="dynamic_resources", action="store_true") p.add_argument("--optimistic-resources", dest="optimistic_resources", action="store_true") p.add_argument("--n-backfill", dest="n_backfill", type=int, default=1) + p.add_argument("--backfill-mem-factor", dest="backfill_mem_factor", type=float, default=1.0, + help="memory over-commit allowed for backfill tasks, as a factor of --mem-limit") p.add_argument("--mem-limit", type=float, default=default_mem, help="in MB") p.add_argument("--cpu-limit", type=float, default=8) @@ -135,6 +137,7 @@ def _args_to_config(ns: argparse.Namespace) -> RunnerConfig: mem_limit=ns.mem_limit, cpu_limit=ns.cpu_limit, n_backfill=ns.n_backfill, + backfill_mem_factor=ns.backfill_mem_factor, update_resources=ns.update_resources, dynamic_resources=ns.dynamic_resources, optimistic_resources=ns.optimistic_resources, diff --git a/MC/workflow_runner/o2dpg_runner/config.py b/MC/workflow_runner/o2dpg_runner/config.py index 5769d0223..3bf9735bb 100644 --- a/MC/workflow_runner/o2dpg_runner/config.py +++ b/MC/workflow_runner/o2dpg_runner/config.py @@ -21,6 +21,7 @@ class RunnerConfig: mem_limit: float = 0.0 # MB; 0 means "auto from psutil" cpu_limit: float = 8.0 n_backfill: int = 1 + backfill_mem_factor: float = 1.0 update_resources: Optional[str] = None dynamic_resources: bool = False optimistic_resources: bool = False diff --git a/MC/workflow_runner/o2dpg_runner/executor.py b/MC/workflow_runner/o2dpg_runner/executor.py index ca2d178ab..fe5cb87c6 100644 --- a/MC/workflow_runner/o2dpg_runner/executor.py +++ b/MC/workflow_runner/o2dpg_runner/executor.py @@ -112,6 +112,7 @@ def __init__( mem_limit=config.mem_limit, procs_parallel_max=config.maxjobs, n_backfill_max=config.n_backfill, + backfill_mem_factor=config.backfill_mem_factor, dynamic_resources=config.dynamic_resources, optimistic_resources=config.optimistic_resources, ) diff --git a/MC/workflow_runner/o2dpg_runner/resources.py b/MC/workflow_runner/o2dpg_runner/resources.py index 4b187d90c..79e604b98 100644 --- a/MC/workflow_runner/o2dpg_runner/resources.py +++ b/MC/workflow_runner/o2dpg_runner/resources.py @@ -182,7 +182,7 @@ def __init__( procs_parallel_max: int = 100, n_backfill_max: int = 1, backfill_cpu_factor: float = 1.5, - backfill_mem_factor: float = 1.5, + backfill_mem_factor: float = 1.0, dynamic_resources: bool = False, optimistic_resources: bool = False, ): @@ -308,12 +308,12 @@ def cpu_free_default(self) -> float: return self.boundaries.cpu_limit - self.cpu_booked def mem_free_default(self) -> float: - return self.boundaries.mem_limit - self.mem_booked + return self.boundaries.mem_limit - self.mem_booked - self.mem_booked_backfill def fits_default(self, res: TaskResources) -> bool: return ( self.cpu_booked + res.cpu_assigned <= self.boundaries.cpu_limit - and self.mem_booked + res.mem_assigned <= self.boundaries.mem_limit + and self.mem_booked + self.mem_booked_backfill + res.mem_assigned <= self.boundaries.mem_limit ) def fits_backfill( diff --git a/MC/workflow_runner/o2dpg_runner/tests/test_resources.py b/MC/workflow_runner/o2dpg_runner/tests/test_resources.py index 2a723fc87..83c58fafc 100644 --- a/MC/workflow_runner/o2dpg_runner/tests/test_resources.py +++ b/MC/workflow_runner/o2dpg_runner/tests/test_resources.py @@ -135,3 +135,30 @@ def test_at_proc_cap(): rm.resources[i].nice_value = rm.nice_default rm.book(i, rm.nice_default) assert rm.at_proc_cap() + + +def test_default_gate_counts_memory_booked_in_backfill(): + """Memory held by a backfill task must count against the default-tier gate.""" + rm = _make_rm(cpu=8, mem=16000) + rm.add_task("a", None, 1, 1, 11000) + rm.add_task("b", None, 1, 1, 11000) + rm.resources[0].nice_value = rm.nice_backfill + rm.book(0, rm.nice_backfill) + assert rm.mem_booked == 0 and rm.mem_booked_backfill == 11000 + assert not rm.fits_default(rm.resources[1]) + assert rm.mem_free_default() == 5000 + + +def test_backfill_mem_factor_default_does_not_overcommit_memory(): + rm = _make_rm(cpu=8, mem=16000) + rm.add_task("a", None, 1, 1, 9000) + rm.add_task("b", None, 1, 1, 9000) + rm.resources[0].nice_value = rm.nice_default + rm.book(0, rm.nice_default) + assert not rm.fits_backfill(rm.resources[1]) + rm15 = _make_rm(cpu=8, mem=16000, backfill_mem_factor=1.5) + rm15.add_task("a", None, 1, 1, 9000) + rm15.add_task("b", None, 1, 1, 9000) + rm15.resources[0].nice_value = rm15.nice_default + rm15.book(0, rm15.nice_default) + assert rm15.fits_backfill(rm15.resources[1]) diff --git a/MC/workflow_runner/o2dpg_schedule_simulator.py b/MC/workflow_runner/o2dpg_schedule_simulator.py index bb9e80819..016dc82b2 100755 --- a/MC/workflow_runner/o2dpg_schedule_simulator.py +++ b/MC/workflow_runner/o2dpg_schedule_simulator.py @@ -261,7 +261,7 @@ def _build_rm( cpu_overrides: Optional[Dict[int, float]] = None, n_backfill_max: int = 0, backfill_cpu_factor: float = 1.5, - backfill_mem_factor: float = 1.5, + backfill_mem_factor: float = 1.0, maxjobs: int = 10_000, ) -> Tuple[ResourceManager, Set[int]]: """Fresh ResourceManager with no backfill tier and unlimited job slots. @@ -374,7 +374,7 @@ def simulate( backfill_model: str = "off", n_backfill: int = 1, backfill_cpu_factor: float = 1.5, - backfill_mem_factor: float = 1.5, + backfill_mem_factor: float = 1.0, backfill_slowdown_factor: float = 1.15, maxjobs: int = 10_000, ) -> SimResult: @@ -802,7 +802,7 @@ def optimize_workers( backfill_model: str = "off", n_backfill: int = 1, backfill_cpu_factor: float = 1.5, - backfill_mem_factor: float = 1.5, + backfill_mem_factor: float = 1.0, backfill_slowdown_factor: float = 1.15, maxjobs: int = 10_000, ) -> Tuple[Dict[str, int], float]: @@ -911,7 +911,7 @@ def build_parser() -> argparse.ArgumentParser: help="Maximum concurrent backfill tasks when backfill simulation is enabled.") p.add_argument("--backfill-cpu-factor", type=float, default=1.5, metavar="X", help="Total CPU oversubscription factor allowed for backfill admission.") - p.add_argument("--backfill-mem-factor", type=float, default=1.5, metavar="X", + p.add_argument("--backfill-mem-factor", type=float, default=1.0, metavar="X", help="Total memory oversubscription factor allowed for backfill admission.") p.add_argument("--backfill-slowdown-factor", type=float, default=1.15, metavar="X", help="Walltime multiplier applied to backfill tasks in "