Coverage for tests/test_prepare_utils.py: 100%
466 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-09-17 02:08 -0700
« prev ^ index » next coverage.py v7.15.4, created at 2026-09-17 02:08 -0700
1# This file is part of ctrl_bps_htcondor.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This software is dual licensed under the GNU General Public License and also
10# under a 3-clause BSD license. Recipients may choose which of these licenses
11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt,
12# respectively. If you choose the GPL option then the following text applies
13# (but note that there is still no warranty even if you opt for BSD instead):
14#
15# This program is free software: you can redistribute it and/or modify
16# it under the terms of the GNU General Public License as published by
17# the Free Software Foundation, either version 3 of the License, or
18# (at your option) any later version.
19#
20# This program is distributed in the hope that it will be useful,
21# but WITHOUT ANY WARRANTY; without even the implied warranty of
22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
23# GNU General Public License for more details.
24#
25# You should have received a copy of the GNU General Public License
26# along with this program. If not, see <https://www.gnu.org/licenses/>.
28"""Unit tests for prepare utility functions."""
30import logging
31import os
32import unittest
33from copy import deepcopy
35from networkx import is_isomorphic
37from lsst.ctrl.bps import (
38 BPS_DEFAULTS,
39 BPS_SEARCH_ORDER,
40 BpsConfig,
41 GenericWorkflow,
42 GenericWorkflowExec,
43 GenericWorkflowFile,
44 GenericWorkflowJob,
45)
46from lsst.ctrl.bps.htcondor import lssthtc, prepare_utils
47from lsst.ctrl.bps.htcondor.htcondor_config import HTC_DEFAULTS_URI
48from lsst.ctrl.bps.tests.gw_test_utils import (
49 make_3_label_workflow,
50 make_3_label_workflow_groups_sort,
51 make_lazy_workflow,
52)
53from lsst.daf.butler import Config
54from lsst.utils.tests import temporaryDirectory
56logger = logging.getLogger("lsst.ctrl.bps.htcondor")
58TESTDIR = os.path.abspath(os.path.dirname(__file__))
61class TranslateJobCmdsTestCase(unittest.TestCase):
62 """Test _translate_job_cmds method."""
64 def setUp(self):
65 self.gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask")
66 self.cached_vals = {
67 "profile": {},
68 "bpsUseShared": True,
69 "memoryLimit": 32768,
70 "bpsMakeCommand": True,
71 "bpsUseHTCEnvironment": True,
72 }
74 def testRetryUnlessNone(self):
75 gwjob = GenericWorkflowJob("retryUnless", "label1", executable=self.gw_exec)
76 gwjob.retry_unless_exit = None
77 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
78 self.assertNotIn("retry_until", htc_commands)
80 def testRetryUnlessInt(self):
81 gwjob = GenericWorkflowJob("retryUnlessInt", "label1", executable=self.gw_exec)
82 gwjob.retry_unless_exit = 3
83 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
84 self.assertEqual(int(htc_commands["retry_until"]), gwjob.retry_unless_exit)
86 def testRetryUnlessList(self):
87 gwjob = GenericWorkflowJob("retryUnlessList", "label1", executable=self.gw_exec)
88 gwjob.retry_unless_exit = [1, 2]
89 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
90 self.assertEqual(htc_commands["retry_until"], "member(ExitCode, {1,2})")
92 def testRetryUnlessBad(self):
93 gwjob = GenericWorkflowJob("retryUnlessBad", "label1", executable=self.gw_exec)
94 gwjob.retry_unless_exit = "1,2,3"
95 with self.assertRaises(ValueError) as cm:
96 _ = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
97 self.assertIn("retryUnlessExit", str(cm.exception))
99 def testEnvironmentBasic(self):
100 gwjob = GenericWorkflowJob("jobEnvironment", "label1", executable=self.gw_exec)
101 gwjob.environment = {"TEST_INT": "1", "TEST_STR": "TWO"}
102 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
103 self.assertEqual(htc_commands["environment"], "TEST_INT='1' TEST_STR='TWO'")
105 def testEnvironmentSpaces(self):
106 gwjob = GenericWorkflowJob("jobEnvironment", "label1", executable=self.gw_exec)
107 gwjob.environment = {"TEST_SPACES": "spacey value"}
108 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
109 self.assertEqual(htc_commands["environment"], "TEST_SPACES='spacey value'")
111 def testEnvironmentSingleQuotes(self):
112 gwjob = GenericWorkflowJob("jobEnvironment", "label1", executable=self.gw_exec)
113 gwjob.environment = {"TEST_SINGLE_QUOTES": "spacey 'quoted' value"}
114 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
115 self.assertEqual(htc_commands["environment"], "TEST_SINGLE_QUOTES='spacey ''quoted'' value'")
117 def testEnvironmentDoubleQuotes(self):
118 gwjob = GenericWorkflowJob("jobEnvironment", "label1", executable=self.gw_exec)
119 gwjob.environment = {"TEST_DOUBLE_QUOTES": 'spacey "double" value'}
120 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
121 self.assertEqual(htc_commands["environment"], """TEST_DOUBLE_QUOTES='spacey ""double"" value'""")
123 def testEnvironmentWithEnvVars(self):
124 gwjob = GenericWorkflowJob("jobEnvironment", "label1", executable=self.gw_exec)
125 gwjob.environment = {"TEST_ENV_VAR": "<ENV:CTRL_BPS_DIR>/tests"}
126 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
127 self.assertEqual(htc_commands["environment"], "TEST_ENV_VAR='$ENV(CTRL_BPS_DIR)/tests'")
129 def testPeriodicRelease(self):
130 gwjob = GenericWorkflowJob("periodicRelease", "label1", executable=self.gw_exec)
131 gwjob.request_memory = 2048
132 gwjob.memory_multiplier = 2
133 gwjob.number_of_retries = 3
134 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
135 release = (
136 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
137 "(HoldReasonCode =?= 12 || (HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
138 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
139 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768)"
140 )
141 self.assertEqual(htc_commands["periodic_release"], release)
143 def testPeriodicRemoveNoRetries(self):
144 gwjob = GenericWorkflowJob("periodicRelease", "label1", executable=self.gw_exec)
145 gwjob.request_memory = 2048
146 gwjob.memory_multiplier = 1
147 gwjob.number_of_retries = 0
148 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
149 remove = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)"
150 self.assertEqual(htc_commands["periodic_remove"], remove)
151 self.assertEqual(htc_commands["max_retries"], 0)
153 def testProfileJobCommands(self):
154 requirement_str = 'Machine == "node01.cluster.local"'
155 gwjob = GenericWorkflowJob("requirements", "label1", executable=self.gw_exec)
156 gwjob.request_memory = 2048
157 gwjob.memory_multiplier = 1
158 gwjob.number_of_retries = 0
159 gwjob.profile = {"requirements": requirement_str}
160 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, None, gwjob)
161 self.assertEqual(htc_commands["requirements"], requirement_str)
163 def testProfileCached(self):
164 requirement_str = 'Machine == "node01.cluster.local"'
165 gwjob = GenericWorkflowJob("requirements", "label1", executable=self.gw_exec)
166 gwjob.request_memory = 2048
167 gwjob.memory_multiplier = 1
168 gwjob.number_of_retries = 0
169 cached_vals = dict(self.cached_vals)
170 cached_vals["profile"] = {"requirements": requirement_str}
171 htc_commands = prepare_utils._translate_job_cmds(cached_vals, None, gwjob)
172 self.assertEqual(htc_commands["requirements"], requirement_str)
174 def testArgumentsReplaceWmsVars(self):
175 gwjob = GenericWorkflowJob("job1", "label1", executable=self.gw_exec)
176 gw = GenericWorkflow("test1")
177 gw.add_job(gwjob)
178 gwjob.request_cpus = 1
179 gwjob.request_memory = 2048
180 gwjob.arguments = "run-qbb repo test.qg --summary /a/b/t/jobs/c/d/job-<WMS:attemptNum>-summary.json"
181 new_arguments = "run-qbb repo test.qg --summary /a/b/t/jobs/c/d/job-$$([NumJobStarts])-summary.json"
182 htc_commands = prepare_utils._translate_job_cmds(self.cached_vals, gw, gwjob)
183 self.assertEqual(htc_commands["arguments"], new_arguments)
186class TranslateCommandLineTestCase(unittest.TestCase):
187 """Test _translate_command_line method."""
189 def setUp(self):
190 self.gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask")
191 self.cached_vals = {"bpsUseShared": True, "bpsMakeCommand": True, "bpsUseHTCEnvironment": True}
193 def _make_job(self, name="job1", executable=None, arguments=None):
194 gwjob = GenericWorkflowJob(name, "label1", executable=executable or self.gw_exec)
195 gw = GenericWorkflow("test1")
196 gw.add_job(gwjob)
197 if arguments is not None:
198 gwjob.arguments = arguments
199 return gw, gwjob
201 def testMakeCommandBasic(self):
202 # Default bpsMakeCommand (True), no executable transfer, no arguments.
203 gw, gwjob = self._make_job()
204 jobcmds = prepare_utils._translate_command_line(self.cached_vals, gw, gwjob)
205 self.assertEqual(jobcmds["getenv"], "True")
206 self.assertEqual(jobcmds["executable"], "/dummy/dir/pipetask")
207 self.assertNotIn("transfer_executable", jobcmds)
208 self.assertNotIn("arguments", jobcmds)
210 def testMakeCommandTransferExecutable(self):
211 gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask", transfer_executable=True)
212 gw, gwjob = self._make_job(executable=gw_exec)
213 jobcmds = prepare_utils._translate_command_line(self.cached_vals, gw, gwjob)
214 self.assertEqual(jobcmds["transfer_executable"], "True")
215 self.assertEqual(jobcmds["executable"], "/dummy/dir/pipetask")
217 def testMakeCommandExecutableEnvVar(self):
218 # Environment placeholders in the executable are converted to HTCondor
219 # env syntax when the executable is not transferred.
220 gw_exec = GenericWorkflowExec("test_exec", "<ENV:CTRL_BPS_DIR>/bin/pipetask")
221 gw, gwjob = self._make_job(executable=gw_exec)
222 jobcmds = prepare_utils._translate_command_line(self.cached_vals, gw, gwjob)
223 self.assertEqual(jobcmds["executable"], "$ENV(CTRL_BPS_DIR)/bin/pipetask")
225 def testMakeCommandArguments(self):
226 # Arguments should have cmd, wms, file, and env placeholders replaced.
227 gw, gwjob = self._make_job(arguments="run <FILE:qg> --attempt <WMS:attemptNum> {opt}")
228 gwjob.cmdvals = {"opt": "X"}
229 gwfile = GenericWorkflowFile("qg", src_uri="/path/to/test.qg", wms_transfer=True, job_shared=True)
230 gw.add_job_inputs(gwjob.name, [gwfile])
231 jobcmds = prepare_utils._translate_command_line(self.cached_vals, gw, gwjob)
232 self.assertEqual(jobcmds["arguments"], "run /path/to/test.qg --attempt $$([NumJobStarts]) X")
234 def testPayloadCommandBasic(self):
235 # bpsMakeCommand False wraps the payloadCommand in /bin/bash -c.
236 gw, gwjob = self._make_job(arguments="run -b repo")
237 cached_vals = {
238 "bpsUseShared": True,
239 "bpsMakeCommand": False,
240 "payloadCommand": "setup; {gwjobCommand}",
241 }
242 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
243 self.assertEqual(jobcmds["executable"], "/bin/bash")
244 self.assertEqual(jobcmds["transfer_executable"], "False")
245 self.assertNotIn("getenv", jobcmds)
246 self.assertEqual(jobcmds["arguments"], "-c 'setup; /dummy/dir/pipetask run -b repo'")
248 def testPayloadCommandNoArguments(self):
249 # No argument to the gwjob command
250 gw, gwjob = self._make_job(arguments="")
251 cached_vals = {
252 "bpsUseShared": True,
253 "bpsMakeCommand": False,
254 "payloadCommand": "setup; {gwjobCommand}",
255 }
256 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
257 self.assertEqual(jobcmds["executable"], "/bin/bash")
258 self.assertEqual(jobcmds["transfer_executable"], "False")
259 self.assertNotIn("getenv", jobcmds)
260 self.assertEqual(jobcmds["arguments"], "-c 'setup; /dummy/dir/pipetask '")
262 def testPayloadCommandStripsNewlines(self):
263 gw, gwjob = self._make_job(arguments="run")
264 cached_vals = {
265 "bpsUseShared": True,
266 "bpsMakeCommand": False,
267 "payloadCommand": "setup;\n{gwjobCommand}",
268 }
269 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
270 self.assertEqual(jobcmds["arguments"], "-c 'setup;/dummy/dir/pipetask run'")
272 def testPayloadCommandExecutableEnvVar(self):
273 # Environment placeholders in the executable use shell syntax here.
274 gw_exec = GenericWorkflowExec("test_exec", "<ENV:CTRL_BPS_DIR>/bin/pipetask")
275 gw, gwjob = self._make_job(executable=gw_exec, arguments="go")
276 cached_vals = {"bpsUseShared": True, "bpsMakeCommand": False, "payloadCommand": "{gwjobCommand}"}
277 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
278 self.assertEqual(jobcmds["arguments"], "-c '${CTRL_BPS_DIR}/bin/pipetask go'")
280 def testPayloadCommandTransferExecutable(self):
281 # Transferred executable is added to the job inputs and chmod'd.
282 gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask", transfer_executable=True)
283 gw, gwjob = self._make_job(executable=gw_exec, arguments="sub")
284 cached_vals = {"bpsUseShared": True, "bpsMakeCommand": False, "payloadCommand": "{gwjobCommand}"}
285 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
286 self.assertEqual(jobcmds["arguments"], "-c 'chmod u+x pipetask; ./pipetask sub'")
287 input_names = [f.name for f in gw.get_job_inputs(gwjob.name, data=True)]
288 self.assertIn("test_exec", input_names)
290 def testEnvironment(self):
291 gw, gwjob = self._make_job()
292 gwjob.environment = {"TEST_INT": 1, "TEST_STR": "TWO"}
293 jobcmds = prepare_utils._translate_command_line(self.cached_vals, gw, gwjob)
294 self.assertEqual(jobcmds["environment"], "TEST_INT='1' TEST_STR='TWO'")
296 def testPayloadCommandEnvironmentShell(self):
297 # Exports in commands, no environment in jobcmds
298 gw_exec = GenericWorkflowExec("test_exec", "<ENV:CTRL_BPS_DIR>/bin/pipetask")
299 gw, gwjob = self._make_job(executable=gw_exec, arguments="go")
300 gwjob.environment = {"TEST_ENV_VAR": "<ENV:CTRL_BPS_DIR>/tests"}
301 cached_vals = {
302 "bpsUseShared": True,
303 "bpsMakeCommand": False,
304 "bpsUseHTCEnvironment": False,
305 "payloadCommand": "{gwjobExports} {gwjobCommand}",
306 }
307 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
308 self.assertIn("export TEST_ENV_VAR='${CTRL_BPS_DIR}/tests';", jobcmds["arguments"])
309 self.assertNotIn("environment", jobcmds)
312class TranslateDagCmdsTestCase(unittest.TestCase):
313 """Test _translate_dag_cmds method."""
315 def setUp(self):
316 self.gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask")
318 def testPriority(self):
319 gwjob = GenericWorkflowJob("priority", "label1", executable=self.gw_exec)
320 gwjob.priority = 100
321 dag_commands = prepare_utils._translate_dag_cmds(gwjob)
322 self.assertEqual(dag_commands["priority"], 100)
325class GroupToSubdagTestCase(unittest.TestCase):
326 """Test _group_to_subdag function."""
328 def testBlocking(self):
329 gw = make_3_label_workflow_groups_sort("test1", True)
330 gwjob = gw.get_job("group_order1_10001")
331 config = BpsConfig(
332 {},
333 search_order=BPS_SEARCH_ORDER,
334 defaults=BPS_DEFAULTS,
335 )
337 htc_job = prepare_utils._group_to_subdag(config, gwjob, "the_prefix")
338 self.assertEqual(len(htc_job.subdag), len(gwjob))
341class GatherSiteValuesTestCase(unittest.TestCase):
342 """Test _gather_site_values function."""
344 def testAllThere(self):
345 config = BpsConfig(
346 {},
347 search_order=BPS_SEARCH_ORDER,
348 defaults=BPS_DEFAULTS,
349 )
350 compute_site = "notThere"
351 results = prepare_utils._gather_site_values(config, compute_site)
352 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"])
354 def testNotSpecified(self):
355 config = BpsConfig(
356 {},
357 search_order=BPS_SEARCH_ORDER,
358 defaults=BPS_DEFAULTS,
359 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
360 )
361 compute_site = "notThere"
362 results = prepare_utils._gather_site_values(config, compute_site)
363 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"])
365 def testAttrsProfile(self):
366 test_values = {
367 "bpsNodeset": "DEVSET",
368 "site": {
369 "mycomputer": {
370 "profile": {
371 "condor": {
372 "requirements": '( TARGET.Nodeset == "{bpsNodeset}" )',
373 "+JobNodeset": "{bpsNodeset}",
374 }
375 }
376 }
377 },
378 }
379 config = BpsConfig(
380 test_values,
381 search_order=BPS_SEARCH_ORDER,
382 defaults=BPS_DEFAULTS,
383 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
384 )
385 results = prepare_utils._gather_site_values(config, "mycomputer")
386 self.assertEqual(results["profile"], {"requirements": '( TARGET.Nodeset == "DEVSET" )'})
387 self.assertEqual(results["attrs"], {"JobNodeset": "DEVSET"})
390class GatherLabelValuesTestCase(unittest.TestCase):
391 """Test _gather_labels_values function."""
393 def testClusterLabel(self):
394 # Test cluster value overrides pipetask.
395 config = BpsConfig(
396 {
397 "cluster": {
398 "label1": {
399 "releaseExpr": "cluster_val",
400 "overwriteJobFiles": False,
401 "profile": {"condor": {"prof_val1": 3}},
402 }
403 },
404 "pipetask": {"label1": {"releaseExpr": "pipetask_val"}},
405 "site": {"site1": {}},
406 },
407 search_order=BPS_SEARCH_ORDER,
408 defaults=BPS_DEFAULTS,
409 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
410 )
411 results = prepare_utils._gather_label_values(config, "label1")
412 self.assertEqual(
413 results,
414 {
415 "attrs": {},
416 "profile": {"prof_val1": 3},
417 "releaseExpr": "cluster_val",
418 "overwriteJobFiles": False,
419 "bpsMakeCommand": True,
420 "bpsUseHTCEnvironment": True,
421 "bpsUseShared": True,
422 "memoryLimit": 491520,
423 },
424 )
426 def testPipetaskLabel(self):
427 label = "label1"
428 config = BpsConfig(
429 {
430 "pipetask": {
431 "label1": {
432 "releaseExpr": "pipetask_val",
433 "overwriteJobFiles": False,
434 "profile": {"condor": {"prof_val1": 3}},
435 }
436 },
437 "site": {"site1": {}},
438 },
439 search_order=BPS_SEARCH_ORDER,
440 defaults=BPS_DEFAULTS,
441 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
442 )
443 results = prepare_utils._gather_label_values(config, label)
444 self.assertEqual(
445 results,
446 {
447 "attrs": {},
448 "bpsMakeCommand": True,
449 "bpsUseHTCEnvironment": True,
450 "bpsUseShared": True,
451 "memoryLimit": 491520,
452 "overwriteJobFiles": False,
453 "profile": {"prof_val1": 3},
454 "releaseExpr": "pipetask_val",
455 },
456 )
458 def testNoSection(self):
459 label = "notThere"
460 config = BpsConfig(
461 {"site": {"site1": {}}},
462 search_order=BPS_SEARCH_ORDER,
463 defaults=BPS_DEFAULTS,
464 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
465 )
466 results = prepare_utils._gather_label_values(config, label)
467 self.assertEqual(
468 results,
469 {
470 "attrs": {},
471 "profile": {},
472 "overwriteJobFiles": True,
473 "bpsMakeCommand": True,
474 "bpsUseHTCEnvironment": True,
475 "bpsUseShared": True,
476 "memoryLimit": 491520,
477 },
478 )
480 def testNoOverwriteSpecified(self):
481 label = "notthere"
482 config = BpsConfig(
483 {"site": {"site1": {}}, "memoryLimit": 491520},
484 search_order=BPS_SEARCH_ORDER,
485 defaults={},
486 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
487 )
488 results = prepare_utils._gather_label_values(config, label)
489 self.assertEqual(
490 results,
491 {
492 "attrs": {},
493 "profile": {},
494 "overwriteJobFiles": True,
495 "bpsMakeCommand": True,
496 "bpsUseHTCEnvironment": True,
497 "bpsUseShared": False,
498 "memoryLimit": 491520,
499 },
500 )
502 def testFinalJob(self):
503 label = "finalJob"
504 config = BpsConfig(
505 {"site": {"site1": {}}, "finalJob": {"profile": {"condor": {"prof_val2": 6, "+attr_val1": 5}}}},
506 search_order=BPS_SEARCH_ORDER,
507 defaults=BPS_DEFAULTS,
508 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
509 )
510 results = prepare_utils._gather_label_values(config, label)
511 self.assertEqual(
512 results,
513 {
514 "attrs": {"attr_val1": 5},
515 "profile": {"prof_val2": 6},
516 "overwriteJobFiles": False,
517 "bpsMakeCommand": True,
518 "bpsUseHTCEnvironment": True,
519 "bpsUseShared": True,
520 "memoryLimit": 491520,
521 },
522 )
524 def testGlobalNodeset(self):
525 config = BpsConfig(
526 {"nodeset": "global_node_set_{campaign}", "campaign": "DRP"},
527 search_order=BPS_SEARCH_ORDER,
528 defaults=BPS_DEFAULTS,
529 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
530 )
531 results = prepare_utils._gather_label_values(config, "label1")
532 self.assertEqual(results["nodeset"], "global_node_set_DRP")
534 def testSiteNodeset(self):
535 config = BpsConfig(
536 {
537 "nodeset": "global_node_set_{campaign}",
538 "campaign": "DRP",
539 "site": {"fr": {"nodeset": "fr_node_set_{campaign}", "siteVar": "frSiteVal"}},
540 "computeSite": "fr",
541 },
542 search_order=BPS_SEARCH_ORDER,
543 defaults=BPS_DEFAULTS,
544 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
545 )
546 results = prepare_utils._gather_label_values(config, "label1")
547 self.assertEqual(results["nodeset"], "fr_node_set_DRP")
548 self.assertEqual(results["siteVar"], "frSiteVal")
550 def testBpsMakeCommandFalse(self):
551 config = BpsConfig(
552 {
553 "bpsMakeCommand": False,
554 },
555 search_order=BPS_SEARCH_ORDER,
556 defaults=BPS_DEFAULTS,
557 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
558 )
559 results = prepare_utils._gather_label_values(config, "label1")
560 self.assertIn("payloadCommand", results)
561 self.assertIn("gwjobCommand", results["payloadCommand"])
564class CreateCheckJobTestCase(unittest.TestCase):
565 """Test _create_check_job function."""
567 def testSuccess(self):
568 group_job_name = "group_order1_val1a"
569 job_label = "order1"
570 job = prepare_utils._create_check_job(group_job_name, job_label, {})
571 self.assertIn(group_job_name, job.name)
572 self.assertEqual(job.label, job_label)
573 self.assertIn("check_group_status.sub", job.subfile)
574 self.assertNotIn("job_nodeset", job.dagcmds["vars"])
576 def testNodeSetSuccess(self):
577 group_job_name = "group_order1_val1a"
578 job_label = "order1"
579 job = prepare_utils._create_check_job(group_job_name, job_label, {"nodeset": "custom_nodeset"})
580 self.assertIn(group_job_name, job.name)
581 self.assertEqual(job.label, job_label)
582 self.assertIn("check_group_status.sub", job.subfile)
583 self.assertIn("job_nodeset", job.dagcmds["vars"])
586class CreatePeriodicReleaseExprTestCase(unittest.TestCase):
587 """Test _create_periodic_release_expr function."""
589 def setUp(self):
590 self.maxDiff = None
592 def testNoReleaseExpr(self):
593 results = prepare_utils._create_periodic_release_expr(2048, 1, 32768, "")
594 self.assertEqual(results, "")
596 def testMultiplierNone(self):
597 results = prepare_utils._create_periodic_release_expr(2048, None, 32768, "")
598 self.assertEqual(results, "")
600 def testJustMemoryReleaseExpr(self):
601 self.maxDiff = None # so test error shows entire strings
602 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "")
603 truth = (
604 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
605 "(HoldReasonCode =?= 12 || "
606 "(HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
607 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
608 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768)"
609 )
610 self.assertEqual(results, truth)
612 def testJustUserReleaseExpr(self):
613 results = prepare_utils._create_periodic_release_expr(2048, 1, 32768, "True")
614 truth = (
615 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
616 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)"
617 )
618 self.assertEqual(results, truth)
620 def testJustUserReleaseExprMultiplierNone(self):
621 results = prepare_utils._create_periodic_release_expr(2048, None, 32768, "True")
622 truth = (
623 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
624 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)"
625 )
626 self.assertEqual(results, truth)
628 def testMemoryAndUserReleaseExpr(self):
629 self.maxDiff = None # so test error shows entire strings
630 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "True")
631 truth = (
632 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
633 "(HoldReasonCode =?= 12 || (HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
634 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
635 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768 || "
636 "HoldReasonCode =!= 1 && True)"
637 )
638 self.assertEqual(results, truth)
641class CreatePeriodicRemoveExprTestCase(unittest.TestCase):
642 """Test _create_periodic_release_expr function."""
644 def testBasicRemoveExpr(self):
645 """Function assumes only called if max_retries >= 0."""
646 results = prepare_utils._create_periodic_remove_expr(2048, 1, 32768)
647 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)"
648 self.assertEqual(results, truth)
650 def testBasicRemoveExprMultiplierNone(self):
651 """Function assumes only called if max_retries >= 0."""
652 results = prepare_utils._create_periodic_remove_expr(2048, None, 32768)
653 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)"
654 self.assertEqual(results, truth)
656 def testMemoryRemoveExpr(self):
657 self.maxDiff = None # so test error shows entire strings
658 results = prepare_utils._create_periodic_remove_expr(2048, 2, 32768)
659 truth = (
660 "JobStatus == 5 && (NumJobStarts > JobMaxRetries || "
661 "((HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
662 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
663 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) == 32768))"
664 )
665 self.assertEqual(results, truth)
668class HandleJobOutputsTestCase(unittest.TestCase):
669 """Test _handle_job_outputs function."""
671 def setUp(self):
672 self.job_name = "test_job"
673 self.out_prefix = "/test/prefix"
675 def tearDown(self):
676 pass
678 def testNoOutputsSharedFilesystem(self):
679 """Test with shared filesystem and no outputs."""
680 mock_workflow = unittest.mock.Mock()
681 mock_workflow.get_job_outputs.return_value = []
683 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix)
685 self.assertEqual(result, {"transfer_output_files": '""'})
687 def testWithOutputsSharedFilesystem(self):
688 """Test with shared filesystem and outputs present (still empty)."""
689 mock_workflow = unittest.mock.Mock()
690 mock_workflow.get_job_outputs.return_value = [
691 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
692 ]
694 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix)
696 self.assertEqual(result, {"transfer_output_files": '""'})
698 def testNoOutputsNoSharedFilesystem(self):
699 """Test without shared filesystem and no outputs."""
700 mock_workflow = unittest.mock.Mock()
701 mock_workflow.get_job_outputs.return_value = []
703 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
705 self.assertEqual(result, {"transfer_output_files": '""'})
707 def testWithAnOutputNoSharedFilesystem(self):
708 """Test without shared filesystem and single output file."""
709 mock_workflow = unittest.mock.Mock()
710 mock_workflow.get_job_outputs.return_value = [
711 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
712 ]
714 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
716 expected = {
717 "transfer_output_files": "output.txt",
718 "transfer_output_remaps": '"output.txt=/path/to/output.txt"',
719 }
720 self.assertEqual(result, expected)
722 def testWithOutputsNoSharedFilesystem(self):
723 """Test without shared filesystem and multiple output files."""
724 mock_workflow = unittest.mock.Mock()
725 mock_workflow.get_job_outputs.return_value = [
726 GenericWorkflowFile(name="output1.txt", src_uri="/path/output1.txt"),
727 GenericWorkflowFile(name="output2.txt", src_uri="/another/path/output2.txt"),
728 ]
730 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
732 expected = {
733 "transfer_output_files": "output1.txt,output2.txt",
734 "transfer_output_remaps": '"output1.txt=/path/output1.txt;output2.txt=/another/path/output2.txt"',
735 }
736 self.assertEqual(result, expected)
738 @unittest.mock.patch("lsst.ctrl.bps.htcondor.prepare_utils._LOG")
739 def testLogging(self, mock_log):
740 mock_workflow = unittest.mock.Mock()
741 mock_workflow.get_job_outputs.return_value = [
742 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
743 ]
745 prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
747 self.assertTrue(mock_log.debug.called)
748 debug_calls = mock_log.debug.call_args_list
749 self.assertTrue(any("src_uri=" in str(call) for call in debug_calls))
750 self.assertTrue(any("transfer_output_files=" in str(call) for call in debug_calls))
751 self.assertTrue(any("transfer_output_remaps=" in str(call) for call in debug_calls))
754class CreateJobTestCase(unittest.TestCase):
755 """Test _create_job function."""
757 def setUp(self):
758 self.generic_workflow = make_3_label_workflow("test1", True)
759 self.template = "{label}/{tract}/{patch}/{band}/{subfilter}/{physical_filter}/{visit}/{exposure}"
761 def testNoOverwrite(self):
762 cached_values = {
763 "bpsUseShared": True,
764 "overwriteJobFiles": False,
765 "memoryLimit": 491520,
766 "profile": {},
767 "attrs": {},
768 }
769 gwjob = self.generic_workflow.get_final()
770 out_prefix = "submit"
771 htc_job = prepare_utils._create_job(
772 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
773 )
774 self.assertEqual(htc_job.name, gwjob.name)
775 self.assertEqual(htc_job.label, gwjob.label)
776 self.assertIn("NumJobStarts", htc_job.cmds["output"])
777 self.assertIn("NumJobStarts", htc_job.cmds["error"])
778 self.assertNotIn("NumJobStarts", htc_job.cmds["log"])
779 self.assertTrue(htc_job.cmds["error"].endswith(".out"))
780 self.assertTrue(htc_job.cmds["output"].endswith(".out"))
781 self.assertTrue(htc_job.cmds["log"].endswith(".log"))
783 def testNodesetWithNoRequirements(self):
784 cached_values = {
785 "bpsUseShared": True,
786 "overwriteJobFiles": False,
787 "memoryLimit": 491520,
788 "profile": {},
789 "attrs": {},
790 "nodeset": "set1",
791 }
792 gwjob = self.generic_workflow.get_job("label1_10002_11")
793 out_prefix = "temp"
794 htc_job = prepare_utils._create_job(
795 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
796 )
797 self.assertEqual(htc_job.cmds["requirements"], '( Target.Nodeset == "set1" )')
798 self.assertEqual(htc_job.attrs["JobNodeset"], "set1")
800 def testNodesetWithRequirements(self):
801 cached_values = {
802 "bpsUseShared": True,
803 "overwriteJobFiles": False,
804 "memoryLimit": 491520,
805 "profile": {"requirements": "dummy_val == 3"},
806 "attrs": {},
807 "nodeset": "set1",
808 }
809 gwjob = self.generic_workflow.get_job("label1_10002_11")
810 out_prefix = "temp"
811 htc_job = prepare_utils._create_job(
812 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
813 )
814 self.assertEqual(htc_job.cmds["requirements"], '(dummy_val == 3) && ( Target.Nodeset == "set1" )')
815 self.assertEqual(htc_job.attrs["JobNodeset"], "set1")
818class ReplaceWmsVarsTestCase(unittest.TestCase):
819 """Test _replace_wms_vars function."""
821 def testNoWmsVar(self):
822 orig_string = "whatever <Other:notThere> whatnot"
823 updated_string = prepare_utils._replace_wms_vars(orig_string)
824 self.assertEqual(orig_string, updated_string)
826 def testAttemptNum(self):
827 orig_string = "whatever <WMS:attemptNum> whatnot"
828 updated_string = prepare_utils._replace_wms_vars(orig_string)
829 self.assertEqual("whatever $$([NumJobStarts]) whatnot", updated_string)
831 def testUnrecognized(self):
832 orig_string = "whatever <WMS:notThere> whatnot"
833 with self.assertLogs(level="INFO") as cm_log:
834 with self.assertRaises(KeyError):
835 _ = prepare_utils._replace_wms_vars(orig_string)
836 self.assertRegex(cm_log.output[0], "Unrecognized WMS placeholder: notThere")
839class UpdateJobSummaryTestCase(unittest.TestCase):
840 """Test _update_job_summary function."""
842 def setUp(self):
843 self.run = "u_testuser_DM-53494_20260220T001651Z"
844 self.filename = f"{self.run}_ctrl.info.json"
845 self.data = {
846 "mycomputer": {
847 "24390.0": {
848 "ClusterId": 24390,
849 "GlobalJobId": "mycomputer#24390.0#1771546612",
850 "bps_run": f"{self.run}_ctrl",
851 "bps_isjob": "True",
852 "bps_payload": "DM-53494",
853 "bps_project": "dev",
854 "bps_runsite": "site1",
855 "bps_campaign": "ci_rc2",
856 "bps_operator": "testuser",
857 "bps_run_quanta": "",
858 "bps_job_summary": "buildQuantumGraph:1;preparePayloadWorkflow:1;dummyJob:1",
859 "bps_wms_service": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorService",
860 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_workflow.HTCondorWorkflow",
861 "bps_wms_config_path": "dagman.conf",
862 }
863 }
864 }
865 self.mapping = f"{self.run}:preparePayloadWorkflow"
866 self.add_summary = "pipetaskInit:1;isr:6;finalJob:1"
868 def testLazyMapping(self):
869 dag_info = deepcopy(self.data)
870 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = self.mapping
871 with temporaryDirectory() as tmp_dir:
872 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info)
874 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
876 _, results = lssthtc.read_dag_info(str(tmp_dir))
877 self.assertEqual(
878 results["mycomputer"]["24390.0"]["bps_job_summary"],
879 f"buildQuantumGraph:1;preparePayloadWorkflow:1;{self.add_summary};dummyJob:1",
880 )
882 def testNoLazyMapping(self):
883 # No bps_lazy_mapping at all
884 with temporaryDirectory() as tmp_dir:
885 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", self.data)
887 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
889 _, results = lssthtc.read_dag_info(str(tmp_dir))
890 self.assertEqual(
891 results["mycomputer"]["24390.0"]["bps_job_summary"],
892 f"{self.data['mycomputer']['24390.0']['bps_job_summary']};{self.add_summary}",
893 )
895 def testNoEntryLazyMapping(self):
896 # bps_lazy_mapping exists, but doesn't include this job
897 dag_info = deepcopy(self.data)
898 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = "other:preparePayloadWorkflow"
899 with temporaryDirectory() as tmp_dir:
900 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info)
902 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
904 _, results = lssthtc.read_dag_info(str(tmp_dir))
905 self.assertEqual(
906 results["mycomputer"]["24390.0"]["bps_job_summary"],
907 f"{self.data['mycomputer']['24390.0']['bps_job_summary']};{self.add_summary}",
908 )
911class ReplaceCmdVarsTestCase(unittest.TestCase):
912 """Test _replace_cmd_vars function."""
914 def testKeyError(self):
915 gwjob = GenericWorkflowJob("job1", "label1")
916 with self.assertLogs(level="DEBUG") as cm_log:
917 with self.assertRaisesRegex(KeyError, ".*notthere.*"):
918 _ = prepare_utils._replace_cmd_vars("{notthere}", gwjob)
919 self.assertRegex(cm_log.output[0], ".*replacement for 'notthere' not provided.*")
922class GenericWorkflowToHTCondorDAG(unittest.TestCase):
923 """Test _generic_workflow_to_htcondor_dag function."""
925 def testRegularWorkflow(self):
926 timestamp = "20260130T211713Z"
927 generic_workflow = make_3_label_workflow("test1", True)
928 config = BpsConfig(
929 {
930 "bpsUseShared": True,
931 "overwriteJobFiles": False,
932 "profile": {"requirements": "dummy_val == 3"},
933 "attrs": {},
934 "nodeset": "set1", # this shouldn't be used with auto-provisioning
935 "provisionResources": True,
936 "provisioning": {"provisioningMaxWallTime": 1200},
937 "bps_defined": {"timestamp": timestamp},
938 "saveHTCdot": True,
939 },
940 defaults=Config(HTC_DEFAULTS_URI),
941 )
943 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir")
944 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items())
945 self.assertIsNotNone(results.graph["final_job"])
946 self.assertTrue(is_isomorphic(results, generic_workflow))
947 self.assertTrue(results.graph["write_dot"])
949 def testLazyWorkflow(self):
950 timestamp = "20260130T211713Z"
951 generic_workflow = make_lazy_workflow("test1", True)
952 config = BpsConfig(
953 {
954 "bpsUseShared": True,
955 "overwriteJobFiles": False,
956 "profile": {"requirements": "dummy_val == 3"},
957 "attrs": {},
958 "nodeset": "set1", # this shouldn't be used with auto-provisioning
959 "provisionResources": True,
960 "provisioning": {"provisioningMaxWallTime": 1200},
961 "bps_defined": {"timestamp": timestamp},
962 },
963 defaults=Config(HTC_DEFAULTS_URI),
964 )
966 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir")
967 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items())
968 self.assertIsNotNone(results.graph["final_job"])
969 # Can't test isomorphic because HTCDag will have additional job for
970 # the lazy dagman job.
971 self.assertTrue(generic_workflow.nodes <= results.nodes)
972 self.assertFalse(results.graph["write_dot"])
975if __name__ == "__main__":
976 unittest.main()