From 5389a2d3ace5c2e8252465fc896966ad928f8056 Mon Sep 17 00:00:00 2001 From: saivishal1999 Date: Thu, 13 Aug 2026 08:57:23 -0700 Subject: [PATCH 1/6] feat(dse): add pin_nodes_to_first_step option to lock DSE steps to same nodes When enabled, all DSE steps after the first run on the same node set as step 1, improving cross-step comparability by eliminating node variance. Co-Authored-By: Claude Sonnet 4.6 (1M context) --- src/cloudai/configurator/cloudai_gym.py | 14 ++++++++++++++ src/cloudai/models/workload.py | 4 ++++ 2 files changed, 18 insertions(+) diff --git a/src/cloudai/configurator/cloudai_gym.py b/src/cloudai/configurator/cloudai_gym.py index cdafd5ce5..2413d3e2c 100644 --- a/src/cloudai/configurator/cloudai_gym.py +++ b/src/cloudai/configurator/cloudai_gym.py @@ -60,6 +60,7 @@ def __init__( self.reward_function = Registry().get_reward_function(test_run.test.agent_reward_function) self.params: EnvParams | None = EnvParams.from_test(test_run.test) self.trajectory = Trajectory(iteration_dir=self.iteration_dir) + self._pinned_nodes: list[str] = [] super().__init__() @property @@ -158,6 +159,8 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: new_tr = copy.deepcopy(self.test_run) new_tr.output_path = self.runner.get_job_output_path(new_tr) + if self.test_run.test.pin_nodes_to_first_step and self._pinned_nodes: + new_tr.nodes = self._pinned_nodes self.runner.test_scenario.test_runs = [new_tr] self.runner.shutting_down = False @@ -169,6 +172,17 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: except Exception as e: logging.error(f"Error running step {self.test_run.step}: {e}") + if self.test_run.test.pin_nodes_to_first_step and not self._pinned_nodes and self.runner.jobs: + job = next(iter(self.runner.jobs.values())) + out, _ = self.runner.system.fetch_command_output( + f"sacct -j {job.id} -p --noheader -X --format=NodeList" + ) + spec = out.splitlines()[0] if out.splitlines() else "" + nodes = spec.strip().replace("|", "") + if nodes and nodes != "Unknown": + self._pinned_nodes = [nodes] + logging.info(f"Pinned DSE nodes to: {nodes}") + if self.runner.test_scenario.test_runs and self.runner.test_scenario.test_runs[0].output_path.exists(): self.test_run = self.runner.test_scenario.test_runs[0] else: diff --git a/src/cloudai/models/workload.py b/src/cloudai/models/workload.py index 22c3c04ad..3c7206f10 100644 --- a/src/cloudai/models/workload.py +++ b/src/cloudai/models/workload.py @@ -123,6 +123,10 @@ class TestDefinition(BaseModel, ABC): agent_metrics: list[str] = Field(default=["default"]) agent_reward_function: str = "inverse" agent_config: dict[str, Any] | None = Field(default=None, description="Agent configuration.") + pin_nodes_to_first_step: bool = Field( + default=False, + description="If True, all DSE steps after the first will be pinned to the same nodes as step 1.", + ) env_params: dict[str, EnvParamSpec] = Field( default_factory=dict, description=( From 4a668e97e6d3f38c344bdb1e088032a4f1547b6c Mon Sep 17 00:00:00 2001 From: saivishal1999 Date: Thu, 13 Aug 2026 09:06:39 -0700 Subject: [PATCH 2/6] fix(dse): fix pin_nodes_to_first_step node capture after run() runner.jobs is cleared on job completion so was empty when we tried to read it. Store last_submitted_job_id on SlurmRunner at submit time and use it to query sacct for the NodeList after step 1 completes. Co-Authored-By: Claude Sonnet 4.6 (1M context) --- src/cloudai/configurator/cloudai_gym.py | 21 +++++++++++---------- src/cloudai/systems/slurm/slurm_runner.py | 2 ++ 2 files changed, 13 insertions(+), 10 deletions(-) diff --git a/src/cloudai/configurator/cloudai_gym.py b/src/cloudai/configurator/cloudai_gym.py index 2413d3e2c..580a96a48 100644 --- a/src/cloudai/configurator/cloudai_gym.py +++ b/src/cloudai/configurator/cloudai_gym.py @@ -172,16 +172,17 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: except Exception as e: logging.error(f"Error running step {self.test_run.step}: {e}") - if self.test_run.test.pin_nodes_to_first_step and not self._pinned_nodes and self.runner.jobs: - job = next(iter(self.runner.jobs.values())) - out, _ = self.runner.system.fetch_command_output( - f"sacct -j {job.id} -p --noheader -X --format=NodeList" - ) - spec = out.splitlines()[0] if out.splitlines() else "" - nodes = spec.strip().replace("|", "") - if nodes and nodes != "Unknown": - self._pinned_nodes = [nodes] - logging.info(f"Pinned DSE nodes to: {nodes}") + if self.test_run.test.pin_nodes_to_first_step and not self._pinned_nodes: + job_id = getattr(self.runner, "last_submitted_job_id", "") + if job_id: + out, _ = self.runner.system.fetch_command_output( + f"sacct -j {job_id} -p --noheader -X --format=NodeList" + ) + spec = out.splitlines()[0] if out.splitlines() else "" + nodes = spec.strip().replace("|", "") + if nodes and nodes not in ("Unknown", ""): + self._pinned_nodes = [nodes] + logging.info(f"Pinned DSE nodes to: {nodes}") if self.runner.test_scenario.test_runs and self.runner.test_scenario.test_runs[0].output_path.exists(): self.test_run = self.runner.test_scenario.test_runs[0] diff --git a/src/cloudai/systems/slurm/slurm_runner.py b/src/cloudai/systems/slurm/slurm_runner.py index dae0cdb29..1942dae8d 100644 --- a/src/cloudai/systems/slurm/slurm_runner.py +++ b/src/cloudai/systems/slurm/slurm_runner.py @@ -42,6 +42,7 @@ def __init__(self, mode: str, system: System, test_scenario: TestScenario, outpu super().__init__(mode, system, test_scenario, output_path) self.system = cast(SlurmSystem, system) self.cmd_shell = CommandShell() + self.last_submitted_job_id: str = "" def get_job_id(self, stdout: str, stderr: str) -> int | None: match = re.search(r"Submitted batch job (\d+)", stdout) @@ -71,6 +72,7 @@ def _submit_test(self, tr: TestRun) -> SlurmJob: message="Failed to retrieve job ID.", ) logging.info(f"Submitted slurm job: {job_id}") + self.last_submitted_job_id = str(job_id) return SlurmJob(tr, id=job_id) def on_job_submit(self, tr: TestRun) -> None: From b6837d8064a97ccc793fb04d54c2d1650736d570 Mon Sep 17 00:00:00 2001 From: saivishal1999 Date: Thu, 13 Aug 2026 09:16:11 -0700 Subject: [PATCH 3/6] fix(dse): fix pin_nodes_to_first_step using one field only runner.jobs is cleared on job completion. Instead of adding a field to the runner, reuse runner.get_job_id() on the stdout file in the output path to find the job ID, then query sacct for the NodeList. Co-Authored-By: Claude Sonnet 4.6 (1M context) --- src/cloudai/configurator/cloudai_gym.py | 22 ++++++++++++---------- 1 file changed, 12 insertions(+), 10 deletions(-) diff --git a/src/cloudai/configurator/cloudai_gym.py b/src/cloudai/configurator/cloudai_gym.py index 580a96a48..41f085929 100644 --- a/src/cloudai/configurator/cloudai_gym.py +++ b/src/cloudai/configurator/cloudai_gym.py @@ -173,16 +173,18 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: logging.error(f"Error running step {self.test_run.step}: {e}") if self.test_run.test.pin_nodes_to_first_step and not self._pinned_nodes: - job_id = getattr(self.runner, "last_submitted_job_id", "") - if job_id: - out, _ = self.runner.system.fetch_command_output( - f"sacct -j {job_id} -p --noheader -X --format=NodeList" - ) - spec = out.splitlines()[0] if out.splitlines() else "" - nodes = spec.strip().replace("|", "") - if nodes and nodes not in ("Unknown", ""): - self._pinned_nodes = [nodes] - logging.info(f"Pinned DSE nodes to: {nodes}") + get_job_id = getattr(self.runner, "get_job_id", None) + for f in new_tr.output_path.rglob("*.stdout"): + job_id = get_job_id(f.read_text(errors="ignore"), "") if get_job_id else None + if job_id: + out, _ = self.runner.system.fetch_command_output( + f"sacct -j {job_id} -p --noheader -X --format=NodeList" + ) + nodes = out.splitlines()[0].strip().replace("|", "") if out.splitlines() else "" + if nodes and nodes != "Unknown": + self._pinned_nodes = [nodes] + logging.info(f"Pinned DSE nodes to: {nodes}") + break if self.runner.test_scenario.test_runs and self.runner.test_scenario.test_runs[0].output_path.exists(): self.test_run = self.runner.test_scenario.test_runs[0] From 266c92e3df2384399086db41e2ef5d64016456c1 Mon Sep 17 00:00:00 2001 From: saivishal1999 Date: Thu, 13 Aug 2026 13:30:33 -0700 Subject: [PATCH 4/6] chore: remove unused last_submitted_job_id from SlurmRunner Replaced by reading job ID from stdout files in output path. Co-Authored-By: Claude Sonnet 4.6 (1M context) --- src/cloudai/systems/slurm/slurm_runner.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/src/cloudai/systems/slurm/slurm_runner.py b/src/cloudai/systems/slurm/slurm_runner.py index 1942dae8d..dae0cdb29 100644 --- a/src/cloudai/systems/slurm/slurm_runner.py +++ b/src/cloudai/systems/slurm/slurm_runner.py @@ -42,7 +42,6 @@ def __init__(self, mode: str, system: System, test_scenario: TestScenario, outpu super().__init__(mode, system, test_scenario, output_path) self.system = cast(SlurmSystem, system) self.cmd_shell = CommandShell() - self.last_submitted_job_id: str = "" def get_job_id(self, stdout: str, stderr: str) -> int | None: match = re.search(r"Submitted batch job (\d+)", stdout) @@ -72,7 +71,6 @@ def _submit_test(self, tr: TestRun) -> SlurmJob: message="Failed to retrieve job ID.", ) logging.info(f"Submitted slurm job: {job_id}") - self.last_submitted_job_id = str(job_id) return SlurmJob(tr, id=job_id) def on_job_submit(self, tr: TestRun) -> None: From 3bfc22db5b067735feb85a572a3ba0d35b5db046 Mon Sep 17 00:00:00 2001 From: saivishal1999 Date: Thu, 13 Aug 2026 13:37:23 -0700 Subject: [PATCH 5/6] fix(dse): fix pyright error for fetch_command_output on base System fetch_command_output is only on SlurmSystem, not System base class. Use getattr to keep the same defensive pattern as get_job_id. Co-Authored-By: Claude Sonnet 4.6 (1M context) --- src/cloudai/configurator/cloudai_gym.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/src/cloudai/configurator/cloudai_gym.py b/src/cloudai/configurator/cloudai_gym.py index 41f085929..ed615073e 100644 --- a/src/cloudai/configurator/cloudai_gym.py +++ b/src/cloudai/configurator/cloudai_gym.py @@ -174,12 +174,11 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: if self.test_run.test.pin_nodes_to_first_step and not self._pinned_nodes: get_job_id = getattr(self.runner, "get_job_id", None) + fetch_cmd = getattr(self.runner.system, "fetch_command_output", None) for f in new_tr.output_path.rglob("*.stdout"): job_id = get_job_id(f.read_text(errors="ignore"), "") if get_job_id else None - if job_id: - out, _ = self.runner.system.fetch_command_output( - f"sacct -j {job_id} -p --noheader -X --format=NodeList" - ) + if job_id and fetch_cmd: + out, _ = fetch_cmd(f"sacct -j {job_id} -p --noheader -X --format=NodeList") nodes = out.splitlines()[0].strip().replace("|", "") if out.splitlines() else "" if nodes and nodes != "Unknown": self._pinned_nodes = [nodes] From 6689e2a39aa596049cfee30ab5b9ce2c49d1159b Mon Sep 17 00:00:00 2001 From: saivishal1999 Date: Thu, 13 Aug 2026 13:43:09 -0700 Subject: [PATCH 6/6] refactor(dse): rename pin_nodes_to_first_step to pin_nodelist Co-Authored-By: Claude Sonnet 4.6 (1M context) --- src/cloudai/configurator/cloudai_gym.py | 4 ++-- src/cloudai/models/workload.py | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/src/cloudai/configurator/cloudai_gym.py b/src/cloudai/configurator/cloudai_gym.py index ed615073e..1d1e256b6 100644 --- a/src/cloudai/configurator/cloudai_gym.py +++ b/src/cloudai/configurator/cloudai_gym.py @@ -159,7 +159,7 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: new_tr = copy.deepcopy(self.test_run) new_tr.output_path = self.runner.get_job_output_path(new_tr) - if self.test_run.test.pin_nodes_to_first_step and self._pinned_nodes: + if self.test_run.test.pin_nodelist and self._pinned_nodes: new_tr.nodes = self._pinned_nodes self.runner.test_scenario.test_runs = [new_tr] @@ -172,7 +172,7 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: except Exception as e: logging.error(f"Error running step {self.test_run.step}: {e}") - if self.test_run.test.pin_nodes_to_first_step and not self._pinned_nodes: + if self.test_run.test.pin_nodelist and not self._pinned_nodes: get_job_id = getattr(self.runner, "get_job_id", None) fetch_cmd = getattr(self.runner.system, "fetch_command_output", None) for f in new_tr.output_path.rglob("*.stdout"): diff --git a/src/cloudai/models/workload.py b/src/cloudai/models/workload.py index 3c7206f10..21d490276 100644 --- a/src/cloudai/models/workload.py +++ b/src/cloudai/models/workload.py @@ -123,7 +123,7 @@ class TestDefinition(BaseModel, ABC): agent_metrics: list[str] = Field(default=["default"]) agent_reward_function: str = "inverse" agent_config: dict[str, Any] | None = Field(default=None, description="Agent configuration.") - pin_nodes_to_first_step: bool = Field( + pin_nodelist: bool = Field( default=False, description="If True, all DSE steps after the first will be pinned to the same nodes as step 1.", )