Coverage for tests/test_prepare_utils.py: 100%
458 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-15 09:02 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-15 09:02 +0000
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="pipetask run")
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 pipetask run'")
248 def testPayloadCommandStripsNewlines(self):
249 gw, gwjob = self._make_job(arguments="run")
250 cached_vals = {
251 "bpsUseShared": True,
252 "bpsMakeCommand": False,
253 "payloadCommand": "setup;\n{gwjobCommand}",
254 }
255 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
256 self.assertEqual(jobcmds["arguments"], "-c 'setup;/dummy/dir/pipetask run'")
258 def testPayloadCommandExecutableEnvVar(self):
259 # Environment placeholders in the executable use shell syntax here.
260 gw_exec = GenericWorkflowExec("test_exec", "<ENV:CTRL_BPS_DIR>/bin/pipetask")
261 gw, gwjob = self._make_job(executable=gw_exec, arguments="go")
262 cached_vals = {"bpsUseShared": True, "bpsMakeCommand": False, "payloadCommand": "{gwjobCommand}"}
263 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
264 self.assertEqual(jobcmds["arguments"], "-c '${CTRL_BPS_DIR}/bin/pipetask go'")
266 def testPayloadCommandTransferExecutable(self):
267 # Transferred executable is added to the job inputs and chmod'd.
268 gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask", transfer_executable=True)
269 gw, gwjob = self._make_job(executable=gw_exec, arguments="sub")
270 cached_vals = {"bpsUseShared": True, "bpsMakeCommand": False, "payloadCommand": "{gwjobCommand}"}
271 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
272 self.assertEqual(jobcmds["arguments"], "-c 'chmod u+x pipetask; ./pipetask sub'")
273 input_names = [f.name for f in gw.get_job_inputs(gwjob.name, data=True)]
274 self.assertIn("test_exec", input_names)
276 def testEnvironment(self):
277 gw, gwjob = self._make_job()
278 gwjob.environment = {"TEST_INT": "1", "TEST_STR": "TWO"}
279 jobcmds = prepare_utils._translate_command_line(self.cached_vals, gw, gwjob)
280 self.assertEqual(jobcmds["environment"], "TEST_INT='1' TEST_STR='TWO'")
282 def testPayloadCommandEnvironmentShell(self):
283 # Exports in commands, no environment in jobcmds
284 gw_exec = GenericWorkflowExec("test_exec", "<ENV:CTRL_BPS_DIR>/bin/pipetask")
285 gw, gwjob = self._make_job(executable=gw_exec, arguments="go")
286 gwjob.environment = {"TEST_ENV_VAR": "<ENV:CTRL_BPS_DIR>/tests"}
287 cached_vals = {
288 "bpsUseShared": True,
289 "bpsMakeCommand": False,
290 "bpsUseHTCEnvironment": False,
291 "payloadCommand": "{gwjobExports} {gwjobCommand}",
292 }
293 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob)
294 self.assertIn("export TEST_ENV_VAR='${CTRL_BPS_DIR}/tests';", jobcmds["arguments"])
295 self.assertNotIn("environment", jobcmds)
298class TranslateDagCmdsTestCase(unittest.TestCase):
299 """Test _translate_dag_cmds method."""
301 def setUp(self):
302 self.gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask")
304 def testPriority(self):
305 gwjob = GenericWorkflowJob("priority", "label1", executable=self.gw_exec)
306 gwjob.priority = 100
307 dag_commands = prepare_utils._translate_dag_cmds(gwjob)
308 self.assertEqual(dag_commands["priority"], 100)
311class GroupToSubdagTestCase(unittest.TestCase):
312 """Test _group_to_subdag function."""
314 def testBlocking(self):
315 gw = make_3_label_workflow_groups_sort("test1", True)
316 gwjob = gw.get_job("group_order1_10001")
317 config = BpsConfig(
318 {},
319 search_order=BPS_SEARCH_ORDER,
320 defaults=BPS_DEFAULTS,
321 )
323 htc_job = prepare_utils._group_to_subdag(config, gwjob, "the_prefix")
324 self.assertEqual(len(htc_job.subdag), len(gwjob))
327class GatherSiteValuesTestCase(unittest.TestCase):
328 """Test _gather_site_values function."""
330 def testAllThere(self):
331 config = BpsConfig(
332 {},
333 search_order=BPS_SEARCH_ORDER,
334 defaults=BPS_DEFAULTS,
335 )
336 compute_site = "notThere"
337 results = prepare_utils._gather_site_values(config, compute_site)
338 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"])
340 def testNotSpecified(self):
341 config = BpsConfig(
342 {},
343 search_order=BPS_SEARCH_ORDER,
344 defaults=BPS_DEFAULTS,
345 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
346 )
347 compute_site = "notThere"
348 results = prepare_utils._gather_site_values(config, compute_site)
349 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"])
351 def testAttrsProfile(self):
352 test_values = {
353 "bpsNodeset": "DEVSET",
354 "site": {
355 "mycomputer": {
356 "profile": {
357 "condor": {
358 "requirements": '( TARGET.Nodeset == "{bpsNodeset}" )',
359 "+JobNodeset": "{bpsNodeset}",
360 }
361 }
362 }
363 },
364 }
365 config = BpsConfig(
366 test_values,
367 search_order=BPS_SEARCH_ORDER,
368 defaults=BPS_DEFAULTS,
369 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
370 )
371 results = prepare_utils._gather_site_values(config, "mycomputer")
372 self.assertEqual(results["profile"], {"requirements": '( TARGET.Nodeset == "DEVSET" )'})
373 self.assertEqual(results["attrs"], {"JobNodeset": "DEVSET"})
376class GatherLabelValuesTestCase(unittest.TestCase):
377 """Test _gather_labels_values function."""
379 def testClusterLabel(self):
380 # Test cluster value overrides pipetask.
381 config = BpsConfig(
382 {
383 "cluster": {
384 "label1": {
385 "releaseExpr": "cluster_val",
386 "overwriteJobFiles": False,
387 "profile": {"condor": {"prof_val1": 3}},
388 }
389 },
390 "pipetask": {"label1": {"releaseExpr": "pipetask_val"}},
391 "site": {"site1": {}},
392 },
393 search_order=BPS_SEARCH_ORDER,
394 defaults=BPS_DEFAULTS,
395 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
396 )
397 results = prepare_utils._gather_label_values(config, "label1")
398 self.assertEqual(
399 results,
400 {
401 "attrs": {},
402 "profile": {"prof_val1": 3},
403 "releaseExpr": "cluster_val",
404 "overwriteJobFiles": False,
405 "bpsMakeCommand": True,
406 "bpsUseHTCEnvironment": True,
407 "bpsUseShared": True,
408 "memoryLimit": 491520,
409 },
410 )
412 def testPipetaskLabel(self):
413 label = "label1"
414 config = BpsConfig(
415 {
416 "pipetask": {
417 "label1": {
418 "releaseExpr": "pipetask_val",
419 "overwriteJobFiles": False,
420 "profile": {"condor": {"prof_val1": 3}},
421 }
422 },
423 "site": {"site1": {}},
424 },
425 search_order=BPS_SEARCH_ORDER,
426 defaults=BPS_DEFAULTS,
427 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
428 )
429 results = prepare_utils._gather_label_values(config, label)
430 self.assertEqual(
431 results,
432 {
433 "attrs": {},
434 "bpsMakeCommand": True,
435 "bpsUseHTCEnvironment": True,
436 "bpsUseShared": True,
437 "memoryLimit": 491520,
438 "overwriteJobFiles": False,
439 "profile": {"prof_val1": 3},
440 "releaseExpr": "pipetask_val",
441 },
442 )
444 def testNoSection(self):
445 label = "notThere"
446 config = BpsConfig(
447 {"site": {"site1": {}}},
448 search_order=BPS_SEARCH_ORDER,
449 defaults=BPS_DEFAULTS,
450 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
451 )
452 results = prepare_utils._gather_label_values(config, label)
453 self.assertEqual(
454 results,
455 {
456 "attrs": {},
457 "profile": {},
458 "overwriteJobFiles": True,
459 "bpsMakeCommand": True,
460 "bpsUseHTCEnvironment": True,
461 "bpsUseShared": True,
462 "memoryLimit": 491520,
463 },
464 )
466 def testNoOverwriteSpecified(self):
467 label = "notthere"
468 config = BpsConfig(
469 {"site": {"site1": {}}, "memoryLimit": 491520},
470 search_order=BPS_SEARCH_ORDER,
471 defaults={},
472 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
473 )
474 results = prepare_utils._gather_label_values(config, label)
475 self.assertEqual(
476 results,
477 {
478 "attrs": {},
479 "profile": {},
480 "overwriteJobFiles": True,
481 "bpsMakeCommand": True,
482 "bpsUseHTCEnvironment": True,
483 "bpsUseShared": False,
484 "memoryLimit": 491520,
485 },
486 )
488 def testFinalJob(self):
489 label = "finalJob"
490 config = BpsConfig(
491 {"site": {"site1": {}}, "finalJob": {"profile": {"condor": {"prof_val2": 6, "+attr_val1": 5}}}},
492 search_order=BPS_SEARCH_ORDER,
493 defaults=BPS_DEFAULTS,
494 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
495 )
496 results = prepare_utils._gather_label_values(config, label)
497 self.assertEqual(
498 results,
499 {
500 "attrs": {"attr_val1": 5},
501 "profile": {"prof_val2": 6},
502 "overwriteJobFiles": False,
503 "bpsMakeCommand": True,
504 "bpsUseHTCEnvironment": True,
505 "bpsUseShared": True,
506 "memoryLimit": 491520,
507 },
508 )
510 def testGlobalNodeset(self):
511 config = BpsConfig(
512 {"nodeset": "global_node_set_{campaign}", "campaign": "DRP"},
513 search_order=BPS_SEARCH_ORDER,
514 defaults=BPS_DEFAULTS,
515 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
516 )
517 results = prepare_utils._gather_label_values(config, "label1")
518 self.assertEqual(results["nodeset"], "global_node_set_DRP")
520 def testSiteNodeset(self):
521 config = BpsConfig(
522 {
523 "nodeset": "global_node_set_{campaign}",
524 "campaign": "DRP",
525 "site": {"fr": {"nodeset": "fr_node_set_{campaign}", "siteVar": "frSiteVal"}},
526 "computeSite": "fr",
527 },
528 search_order=BPS_SEARCH_ORDER,
529 defaults=BPS_DEFAULTS,
530 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
531 )
532 results = prepare_utils._gather_label_values(config, "label1")
533 self.assertEqual(results["nodeset"], "fr_node_set_DRP")
534 self.assertEqual(results["siteVar"], "frSiteVal")
536 def testBpsMakeCommandFalse(self):
537 config = BpsConfig(
538 {
539 "bpsMakeCommand": False,
540 },
541 search_order=BPS_SEARCH_ORDER,
542 defaults=BPS_DEFAULTS,
543 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
544 )
545 results = prepare_utils._gather_label_values(config, "label1")
546 self.assertIn("payloadCommand", results)
547 self.assertIn("gwjobCommand", results["payloadCommand"])
550class CreateCheckJobTestCase(unittest.TestCase):
551 """Test _create_check_job function."""
553 def testSuccess(self):
554 group_job_name = "group_order1_val1a"
555 job_label = "order1"
556 job = prepare_utils._create_check_job(group_job_name, job_label, {})
557 self.assertIn(group_job_name, job.name)
558 self.assertEqual(job.label, job_label)
559 self.assertIn("check_group_status.sub", job.subfile)
560 self.assertNotIn("job_nodeset", job.dagcmds["vars"])
562 def testNodeSetSuccess(self):
563 group_job_name = "group_order1_val1a"
564 job_label = "order1"
565 job = prepare_utils._create_check_job(group_job_name, job_label, {"nodeset": "custom_nodeset"})
566 self.assertIn(group_job_name, job.name)
567 self.assertEqual(job.label, job_label)
568 self.assertIn("check_group_status.sub", job.subfile)
569 self.assertIn("job_nodeset", job.dagcmds["vars"])
572class CreatePeriodicReleaseExprTestCase(unittest.TestCase):
573 """Test _create_periodic_release_expr function."""
575 def setUp(self):
576 self.maxDiff = None
578 def testNoReleaseExpr(self):
579 results = prepare_utils._create_periodic_release_expr(2048, 1, 32768, "")
580 self.assertEqual(results, "")
582 def testMultiplierNone(self):
583 results = prepare_utils._create_periodic_release_expr(2048, None, 32768, "")
584 self.assertEqual(results, "")
586 def testJustMemoryReleaseExpr(self):
587 self.maxDiff = None # so test error shows entire strings
588 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "")
589 truth = (
590 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
591 "(HoldReasonCode =?= 12 || "
592 "(HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
593 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
594 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768)"
595 )
596 self.assertEqual(results, truth)
598 def testJustUserReleaseExpr(self):
599 results = prepare_utils._create_periodic_release_expr(2048, 1, 32768, "True")
600 truth = (
601 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
602 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)"
603 )
604 self.assertEqual(results, truth)
606 def testJustUserReleaseExprMultiplierNone(self):
607 results = prepare_utils._create_periodic_release_expr(2048, None, 32768, "True")
608 truth = (
609 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
610 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)"
611 )
612 self.assertEqual(results, truth)
614 def testMemoryAndUserReleaseExpr(self):
615 self.maxDiff = None # so test error shows entire strings
616 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "True")
617 truth = (
618 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
619 "(HoldReasonCode =?= 12 || (HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
620 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
621 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768 || "
622 "HoldReasonCode =!= 1 && True)"
623 )
624 self.assertEqual(results, truth)
627class CreatePeriodicRemoveExprTestCase(unittest.TestCase):
628 """Test _create_periodic_release_expr function."""
630 def testBasicRemoveExpr(self):
631 """Function assumes only called if max_retries >= 0."""
632 results = prepare_utils._create_periodic_remove_expr(2048, 1, 32768)
633 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)"
634 self.assertEqual(results, truth)
636 def testBasicRemoveExprMultiplierNone(self):
637 """Function assumes only called if max_retries >= 0."""
638 results = prepare_utils._create_periodic_remove_expr(2048, None, 32768)
639 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)"
640 self.assertEqual(results, truth)
642 def testMemoryRemoveExpr(self):
643 self.maxDiff = None # so test error shows entire strings
644 results = prepare_utils._create_periodic_remove_expr(2048, 2, 32768)
645 truth = (
646 "JobStatus == 5 && (NumJobStarts > JobMaxRetries || "
647 "((HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
648 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
649 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) == 32768))"
650 )
651 self.assertEqual(results, truth)
654class HandleJobOutputsTestCase(unittest.TestCase):
655 """Test _handle_job_outputs function."""
657 def setUp(self):
658 self.job_name = "test_job"
659 self.out_prefix = "/test/prefix"
661 def tearDown(self):
662 pass
664 def testNoOutputsSharedFilesystem(self):
665 """Test with shared filesystem and no outputs."""
666 mock_workflow = unittest.mock.Mock()
667 mock_workflow.get_job_outputs.return_value = []
669 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix)
671 self.assertEqual(result, {"transfer_output_files": '""'})
673 def testWithOutputsSharedFilesystem(self):
674 """Test with shared filesystem and outputs present (still empty)."""
675 mock_workflow = unittest.mock.Mock()
676 mock_workflow.get_job_outputs.return_value = [
677 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
678 ]
680 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix)
682 self.assertEqual(result, {"transfer_output_files": '""'})
684 def testNoOutputsNoSharedFilesystem(self):
685 """Test without shared filesystem and no outputs."""
686 mock_workflow = unittest.mock.Mock()
687 mock_workflow.get_job_outputs.return_value = []
689 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
691 self.assertEqual(result, {"transfer_output_files": '""'})
693 def testWithAnOutputNoSharedFilesystem(self):
694 """Test without shared filesystem and single output file."""
695 mock_workflow = unittest.mock.Mock()
696 mock_workflow.get_job_outputs.return_value = [
697 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
698 ]
700 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
702 expected = {
703 "transfer_output_files": "output.txt",
704 "transfer_output_remaps": '"output.txt=/path/to/output.txt"',
705 }
706 self.assertEqual(result, expected)
708 def testWithOutputsNoSharedFilesystem(self):
709 """Test without shared filesystem and multiple output files."""
710 mock_workflow = unittest.mock.Mock()
711 mock_workflow.get_job_outputs.return_value = [
712 GenericWorkflowFile(name="output1.txt", src_uri="/path/output1.txt"),
713 GenericWorkflowFile(name="output2.txt", src_uri="/another/path/output2.txt"),
714 ]
716 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
718 expected = {
719 "transfer_output_files": "output1.txt,output2.txt",
720 "transfer_output_remaps": '"output1.txt=/path/output1.txt;output2.txt=/another/path/output2.txt"',
721 }
722 self.assertEqual(result, expected)
724 @unittest.mock.patch("lsst.ctrl.bps.htcondor.prepare_utils._LOG")
725 def testLogging(self, mock_log):
726 mock_workflow = unittest.mock.Mock()
727 mock_workflow.get_job_outputs.return_value = [
728 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
729 ]
731 prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
733 self.assertTrue(mock_log.debug.called)
734 debug_calls = mock_log.debug.call_args_list
735 self.assertTrue(any("src_uri=" in str(call) for call in debug_calls))
736 self.assertTrue(any("transfer_output_files=" in str(call) for call in debug_calls))
737 self.assertTrue(any("transfer_output_remaps=" in str(call) for call in debug_calls))
740class CreateJobTestCase(unittest.TestCase):
741 """Test _create_job function."""
743 def setUp(self):
744 self.generic_workflow = make_3_label_workflow("test1", True)
745 self.template = "{label}/{tract}/{patch}/{band}/{subfilter}/{physical_filter}/{visit}/{exposure}"
747 def testNoOverwrite(self):
748 cached_values = {
749 "bpsUseShared": True,
750 "overwriteJobFiles": False,
751 "memoryLimit": 491520,
752 "profile": {},
753 "attrs": {},
754 }
755 gwjob = self.generic_workflow.get_final()
756 out_prefix = "submit"
757 htc_job = prepare_utils._create_job(
758 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
759 )
760 self.assertEqual(htc_job.name, gwjob.name)
761 self.assertEqual(htc_job.label, gwjob.label)
762 self.assertIn("NumJobStarts", htc_job.cmds["output"])
763 self.assertIn("NumJobStarts", htc_job.cmds["error"])
764 self.assertNotIn("NumJobStarts", htc_job.cmds["log"])
765 self.assertTrue(htc_job.cmds["error"].endswith(".out"))
766 self.assertTrue(htc_job.cmds["output"].endswith(".out"))
767 self.assertTrue(htc_job.cmds["log"].endswith(".log"))
769 def testNodesetWithNoRequirements(self):
770 cached_values = {
771 "bpsUseShared": True,
772 "overwriteJobFiles": False,
773 "memoryLimit": 491520,
774 "profile": {},
775 "attrs": {},
776 "nodeset": "set1",
777 }
778 gwjob = self.generic_workflow.get_job("label1_10002_11")
779 out_prefix = "temp"
780 htc_job = prepare_utils._create_job(
781 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
782 )
783 self.assertEqual(htc_job.cmds["requirements"], '( Target.Nodeset == "set1" )')
784 self.assertEqual(htc_job.attrs["JobNodeset"], "set1")
786 def testNodesetWithRequirements(self):
787 cached_values = {
788 "bpsUseShared": True,
789 "overwriteJobFiles": False,
790 "memoryLimit": 491520,
791 "profile": {"requirements": "dummy_val == 3"},
792 "attrs": {},
793 "nodeset": "set1",
794 }
795 gwjob = self.generic_workflow.get_job("label1_10002_11")
796 out_prefix = "temp"
797 htc_job = prepare_utils._create_job(
798 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
799 )
800 self.assertEqual(htc_job.cmds["requirements"], '(dummy_val == 3) && ( Target.Nodeset == "set1" )')
801 self.assertEqual(htc_job.attrs["JobNodeset"], "set1")
804class ReplaceWmsVarsTestCase(unittest.TestCase):
805 """Test _replace_wms_vars function."""
807 def testNoWmsVar(self):
808 orig_string = "whatever <Other:notThere> whatnot"
809 updated_string = prepare_utils._replace_wms_vars(orig_string)
810 self.assertEqual(orig_string, updated_string)
812 def testAttemptNum(self):
813 orig_string = "whatever <WMS:attemptNum> whatnot"
814 updated_string = prepare_utils._replace_wms_vars(orig_string)
815 self.assertEqual("whatever $$([NumJobStarts]) whatnot", updated_string)
817 def testUnrecognized(self):
818 orig_string = "whatever <WMS:notThere> whatnot"
819 with self.assertLogs(level="INFO") as cm_log:
820 with self.assertRaises(KeyError):
821 _ = prepare_utils._replace_wms_vars(orig_string)
822 self.assertRegex(cm_log.output[0], "Unrecognized WMS placeholder: notThere")
825class UpdateJobSummaryTestCase(unittest.TestCase):
826 """Test _update_job_summary function."""
828 def setUp(self):
829 self.run = "u_testuser_DM-53494_20260220T001651Z"
830 self.filename = f"{self.run}_ctrl.info.json"
831 self.data = {
832 "mycomputer": {
833 "24390.0": {
834 "ClusterId": 24390,
835 "GlobalJobId": "mycomputer#24390.0#1771546612",
836 "bps_run": f"{self.run}_ctrl",
837 "bps_isjob": "True",
838 "bps_payload": "DM-53494",
839 "bps_project": "dev",
840 "bps_runsite": "site1",
841 "bps_campaign": "ci_rc2",
842 "bps_operator": "testuser",
843 "bps_run_quanta": "",
844 "bps_job_summary": "buildQuantumGraph:1;preparePayloadWorkflow:1;dummyJob:1",
845 "bps_wms_service": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorService",
846 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_workflow.HTCondorWorkflow",
847 "bps_wms_config_path": "dagman.conf",
848 }
849 }
850 }
851 self.mapping = f"{self.run}:preparePayloadWorkflow"
852 self.add_summary = "pipetaskInit:1;isr:6;finalJob:1"
854 def testLazyMapping(self):
855 dag_info = deepcopy(self.data)
856 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = self.mapping
857 with temporaryDirectory() as tmp_dir:
858 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info)
860 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
862 _, results = lssthtc.read_dag_info(str(tmp_dir))
863 self.assertEqual(
864 results["mycomputer"]["24390.0"]["bps_job_summary"],
865 f"buildQuantumGraph:1;preparePayloadWorkflow:1;{self.add_summary};dummyJob:1",
866 )
868 def testNoLazyMapping(self):
869 # No bps_lazy_mapping at all
870 with temporaryDirectory() as tmp_dir:
871 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", self.data)
873 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
875 _, results = lssthtc.read_dag_info(str(tmp_dir))
876 self.assertEqual(
877 results["mycomputer"]["24390.0"]["bps_job_summary"],
878 f"{self.data['mycomputer']['24390.0']['bps_job_summary']};{self.add_summary}",
879 )
881 def testNoEntryLazyMapping(self):
882 # bps_lazy_mapping exists, but doesn't include this job
883 dag_info = deepcopy(self.data)
884 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = "other:preparePayloadWorkflow"
885 with temporaryDirectory() as tmp_dir:
886 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info)
888 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
890 _, results = lssthtc.read_dag_info(str(tmp_dir))
891 self.assertEqual(
892 results["mycomputer"]["24390.0"]["bps_job_summary"],
893 f"{self.data['mycomputer']['24390.0']['bps_job_summary']};{self.add_summary}",
894 )
897class ReplaceCmdVarsTestCase(unittest.TestCase):
898 """Test _replace_cmd_vars function."""
900 def testKeyError(self):
901 gwjob = GenericWorkflowJob("job1", "label1")
902 with self.assertLogs(level="DEBUG") as cm_log:
903 with self.assertRaisesRegex(KeyError, ".*notthere.*"):
904 _ = prepare_utils._replace_cmd_vars("{notthere}", gwjob)
905 self.assertRegex(cm_log.output[0], ".*replacement for 'notthere' not provided.*")
908class GenericWorkflowToHTCondorDAG(unittest.TestCase):
909 """Test _generic_workflow_to_htcondor_dag function."""
911 def testRegularWorkflow(self):
912 timestamp = "20260130T211713Z"
913 generic_workflow = make_3_label_workflow("test1", True)
914 config = BpsConfig(
915 {
916 "bpsUseShared": True,
917 "overwriteJobFiles": False,
918 "profile": {"requirements": "dummy_val == 3"},
919 "attrs": {},
920 "nodeset": "set1", # this shouldn't be used with auto-provisioning
921 "provisionResources": True,
922 "provisioning": {"provisioningMaxWallTime": 1200},
923 "bps_defined": {"timestamp": timestamp},
924 "saveHTCdot": True,
925 },
926 defaults=Config(HTC_DEFAULTS_URI),
927 )
929 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir")
930 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items())
931 self.assertIsNotNone(results.graph["final_job"])
932 self.assertTrue(is_isomorphic(results, generic_workflow))
933 self.assertTrue(results.graph["write_dot"])
935 def testLazyWorkflow(self):
936 timestamp = "20260130T211713Z"
937 generic_workflow = make_lazy_workflow("test1", True)
938 config = BpsConfig(
939 {
940 "bpsUseShared": True,
941 "overwriteJobFiles": False,
942 "profile": {"requirements": "dummy_val == 3"},
943 "attrs": {},
944 "nodeset": "set1", # this shouldn't be used with auto-provisioning
945 "provisionResources": True,
946 "provisioning": {"provisioningMaxWallTime": 1200},
947 "bps_defined": {"timestamp": timestamp},
948 },
949 defaults=Config(HTC_DEFAULTS_URI),
950 )
952 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir")
953 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items())
954 self.assertIsNotNone(results.graph["final_job"])
955 # Can't test isomorphic because HTCDag will have additional job for
956 # the lazy dagman job.
957 self.assertTrue(generic_workflow.nodes <= results.nodes)
958 self.assertFalse(results.graph["write_dot"])
961if __name__ == "__main__":
962 unittest.main()