Coverage for tests/test_prepare_utils.py: 100%
450 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-29 02:16 -0700
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-29 02:16 -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='${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'")
283class TranslateDagCmdsTestCase(unittest.TestCase):
284 """Test _translate_dag_cmds method."""
286 def setUp(self):
287 self.gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask")
289 def testPriority(self):
290 gwjob = GenericWorkflowJob("priority", "label1", executable=self.gw_exec)
291 gwjob.priority = 100
292 dag_commands = prepare_utils._translate_dag_cmds(gwjob)
293 self.assertEqual(dag_commands["priority"], 100)
296class GroupToSubdagTestCase(unittest.TestCase):
297 """Test _group_to_subdag function."""
299 def testBlocking(self):
300 gw = make_3_label_workflow_groups_sort("test1", True)
301 gwjob = gw.get_job("group_order1_10001")
302 config = BpsConfig(
303 {},
304 search_order=BPS_SEARCH_ORDER,
305 defaults=BPS_DEFAULTS,
306 )
308 htc_job = prepare_utils._group_to_subdag(config, gwjob, "the_prefix")
309 self.assertEqual(len(htc_job.subdag), len(gwjob))
312class GatherSiteValuesTestCase(unittest.TestCase):
313 """Test _gather_site_values function."""
315 def testAllThere(self):
316 config = BpsConfig(
317 {},
318 search_order=BPS_SEARCH_ORDER,
319 defaults=BPS_DEFAULTS,
320 )
321 compute_site = "notThere"
322 results = prepare_utils._gather_site_values(config, compute_site)
323 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"])
325 def testNotSpecified(self):
326 config = BpsConfig(
327 {},
328 search_order=BPS_SEARCH_ORDER,
329 defaults=BPS_DEFAULTS,
330 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
331 )
332 compute_site = "notThere"
333 results = prepare_utils._gather_site_values(config, compute_site)
334 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"])
336 def testAttrsProfile(self):
337 test_values = {
338 "bpsNodeset": "DEVSET",
339 "site": {
340 "mycomputer": {
341 "profile": {
342 "condor": {
343 "requirements": '( TARGET.Nodeset == "{bpsNodeset}" )',
344 "+JobNodeset": "{bpsNodeset}",
345 }
346 }
347 }
348 },
349 }
350 config = BpsConfig(
351 test_values,
352 search_order=BPS_SEARCH_ORDER,
353 defaults=BPS_DEFAULTS,
354 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
355 )
356 results = prepare_utils._gather_site_values(config, "mycomputer")
357 self.assertEqual(results["profile"], {"requirements": '( TARGET.Nodeset == "DEVSET" )'})
358 self.assertEqual(results["attrs"], {"JobNodeset": "DEVSET"})
361class GatherLabelValuesTestCase(unittest.TestCase):
362 """Test _gather_labels_values function."""
364 def testClusterLabel(self):
365 # Test cluster value overrides pipetask.
366 config = BpsConfig(
367 {
368 "cluster": {
369 "label1": {
370 "releaseExpr": "cluster_val",
371 "overwriteJobFiles": False,
372 "profile": {"condor": {"prof_val1": 3}},
373 }
374 },
375 "pipetask": {"label1": {"releaseExpr": "pipetask_val"}},
376 "site": {"site1": {}},
377 },
378 search_order=BPS_SEARCH_ORDER,
379 defaults=BPS_DEFAULTS,
380 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
381 )
382 results = prepare_utils._gather_label_values(config, "label1")
383 self.assertEqual(
384 results,
385 {
386 "attrs": {},
387 "profile": {"prof_val1": 3},
388 "releaseExpr": "cluster_val",
389 "overwriteJobFiles": False,
390 "bpsMakeCommand": True,
391 "bpsUseHTCEnvironment": True,
392 "bpsUseShared": True,
393 "memoryLimit": 491520,
394 },
395 )
397 def testPipetaskLabel(self):
398 label = "label1"
399 config = BpsConfig(
400 {
401 "pipetask": {
402 "label1": {
403 "releaseExpr": "pipetask_val",
404 "overwriteJobFiles": False,
405 "profile": {"condor": {"prof_val1": 3}},
406 }
407 },
408 "site": {"site1": {}},
409 },
410 search_order=BPS_SEARCH_ORDER,
411 defaults=BPS_DEFAULTS,
412 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
413 )
414 results = prepare_utils._gather_label_values(config, label)
415 self.assertEqual(
416 results,
417 {
418 "attrs": {},
419 "bpsMakeCommand": True,
420 "bpsUseHTCEnvironment": True,
421 "bpsUseShared": True,
422 "memoryLimit": 491520,
423 "overwriteJobFiles": False,
424 "profile": {"prof_val1": 3},
425 "releaseExpr": "pipetask_val",
426 },
427 )
429 def testNoSection(self):
430 label = "notThere"
431 config = BpsConfig(
432 {"site": {"site1": {}}},
433 search_order=BPS_SEARCH_ORDER,
434 defaults=BPS_DEFAULTS,
435 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
436 )
437 results = prepare_utils._gather_label_values(config, label)
438 self.assertEqual(
439 results,
440 {
441 "attrs": {},
442 "profile": {},
443 "overwriteJobFiles": True,
444 "bpsMakeCommand": True,
445 "bpsUseHTCEnvironment": True,
446 "bpsUseShared": True,
447 "memoryLimit": 491520,
448 },
449 )
451 def testNoOverwriteSpecified(self):
452 label = "notthere"
453 config = BpsConfig(
454 {"site": {"site1": {}}, "memoryLimit": 491520},
455 search_order=BPS_SEARCH_ORDER,
456 defaults={},
457 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
458 )
459 results = prepare_utils._gather_label_values(config, label)
460 self.assertEqual(
461 results,
462 {
463 "attrs": {},
464 "profile": {},
465 "overwriteJobFiles": True,
466 "bpsMakeCommand": True,
467 "bpsUseHTCEnvironment": True,
468 "bpsUseShared": False,
469 "memoryLimit": 491520,
470 },
471 )
473 def testFinalJob(self):
474 label = "finalJob"
475 config = BpsConfig(
476 {"site": {"site1": {}}, "finalJob": {"profile": {"condor": {"prof_val2": 6, "+attr_val1": 5}}}},
477 search_order=BPS_SEARCH_ORDER,
478 defaults=BPS_DEFAULTS,
479 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
480 )
481 results = prepare_utils._gather_label_values(config, label)
482 self.assertEqual(
483 results,
484 {
485 "attrs": {"attr_val1": 5},
486 "profile": {"prof_val2": 6},
487 "overwriteJobFiles": False,
488 "bpsMakeCommand": True,
489 "bpsUseHTCEnvironment": True,
490 "bpsUseShared": True,
491 "memoryLimit": 491520,
492 },
493 )
495 def testGlobalNodeset(self):
496 config = BpsConfig(
497 {"nodeset": "global_node_set_{campaign}", "campaign": "DRP"},
498 search_order=BPS_SEARCH_ORDER,
499 defaults=BPS_DEFAULTS,
500 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
501 )
502 results = prepare_utils._gather_label_values(config, "label1")
503 self.assertEqual(results["nodeset"], "global_node_set_DRP")
505 def testSiteNodeset(self):
506 config = BpsConfig(
507 {
508 "nodeset": "global_node_set_{campaign}",
509 "campaign": "DRP",
510 "site": {"fr": {"nodeset": "fr_node_set_{campaign}", "siteVar": "frSiteVal"}},
511 "computeSite": "fr",
512 },
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"], "fr_node_set_DRP")
519 self.assertEqual(results["siteVar"], "frSiteVal")
521 def testBpsMakeCommandFalse(self):
522 config = BpsConfig(
523 {
524 "bpsMakeCommand": False,
525 },
526 search_order=BPS_SEARCH_ORDER,
527 defaults=BPS_DEFAULTS,
528 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService",
529 )
530 results = prepare_utils._gather_label_values(config, "label1")
531 self.assertIn("payloadCommand", results)
532 self.assertIn("gwjobCommand", results["payloadCommand"])
535class CreateCheckJobTestCase(unittest.TestCase):
536 """Test _create_check_job function."""
538 def testSuccess(self):
539 group_job_name = "group_order1_val1a"
540 job_label = "order1"
541 job = prepare_utils._create_check_job(group_job_name, job_label, {})
542 self.assertIn(group_job_name, job.name)
543 self.assertEqual(job.label, job_label)
544 self.assertIn("check_group_status.sub", job.subfile)
545 self.assertNotIn("job_nodeset", job.dagcmds["vars"])
547 def testNodeSetSuccess(self):
548 group_job_name = "group_order1_val1a"
549 job_label = "order1"
550 job = prepare_utils._create_check_job(group_job_name, job_label, {"nodeset": "custom_nodeset"})
551 self.assertIn(group_job_name, job.name)
552 self.assertEqual(job.label, job_label)
553 self.assertIn("check_group_status.sub", job.subfile)
554 self.assertIn("job_nodeset", job.dagcmds["vars"])
557class CreatePeriodicReleaseExprTestCase(unittest.TestCase):
558 """Test _create_periodic_release_expr function."""
560 def setUp(self):
561 self.maxDiff = None
563 def testNoReleaseExpr(self):
564 results = prepare_utils._create_periodic_release_expr(2048, 1, 32768, "")
565 self.assertEqual(results, "")
567 def testMultiplierNone(self):
568 results = prepare_utils._create_periodic_release_expr(2048, None, 32768, "")
569 self.assertEqual(results, "")
571 def testJustMemoryReleaseExpr(self):
572 self.maxDiff = None # so test error shows entire strings
573 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "")
574 truth = (
575 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
576 "(HoldReasonCode =?= 12 || "
577 "(HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
578 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
579 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768)"
580 )
581 self.assertEqual(results, truth)
583 def testJustUserReleaseExpr(self):
584 results = prepare_utils._create_periodic_release_expr(2048, 1, 32768, "True")
585 truth = (
586 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
587 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)"
588 )
589 self.assertEqual(results, truth)
591 def testJustUserReleaseExprMultiplierNone(self):
592 results = prepare_utils._create_periodic_release_expr(2048, None, 32768, "True")
593 truth = (
594 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
595 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)"
596 )
597 self.assertEqual(results, truth)
599 def testMemoryAndUserReleaseExpr(self):
600 self.maxDiff = None # so test error shows entire strings
601 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "True")
602 truth = (
603 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && "
604 "(HoldReasonCode =?= 12 || (HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
605 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
606 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768 || "
607 "HoldReasonCode =!= 1 && True)"
608 )
609 self.assertEqual(results, truth)
612class CreatePeriodicRemoveExprTestCase(unittest.TestCase):
613 """Test _create_periodic_release_expr function."""
615 def testBasicRemoveExpr(self):
616 """Function assumes only called if max_retries >= 0."""
617 results = prepare_utils._create_periodic_remove_expr(2048, 1, 32768)
618 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)"
619 self.assertEqual(results, truth)
621 def testBasicRemoveExprMultiplierNone(self):
622 """Function assumes only called if max_retries >= 0."""
623 results = prepare_utils._create_periodic_remove_expr(2048, None, 32768)
624 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)"
625 self.assertEqual(results, truth)
627 def testMemoryRemoveExpr(self):
628 self.maxDiff = None # so test error shows entire strings
629 results = prepare_utils._create_periodic_remove_expr(2048, 2, 32768)
630 truth = (
631 "JobStatus == 5 && (NumJobStarts > JobMaxRetries || "
632 "((HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || "
633 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && "
634 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) == 32768))"
635 )
636 self.assertEqual(results, truth)
639class HandleJobOutputsTestCase(unittest.TestCase):
640 """Test _handle_job_outputs function."""
642 def setUp(self):
643 self.job_name = "test_job"
644 self.out_prefix = "/test/prefix"
646 def tearDown(self):
647 pass
649 def testNoOutputsSharedFilesystem(self):
650 """Test with shared filesystem and no outputs."""
651 mock_workflow = unittest.mock.Mock()
652 mock_workflow.get_job_outputs.return_value = []
654 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix)
656 self.assertEqual(result, {"transfer_output_files": '""'})
658 def testWithOutputsSharedFilesystem(self):
659 """Test with shared filesystem and outputs present (still empty)."""
660 mock_workflow = unittest.mock.Mock()
661 mock_workflow.get_job_outputs.return_value = [
662 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
663 ]
665 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix)
667 self.assertEqual(result, {"transfer_output_files": '""'})
669 def testNoOutputsNoSharedFilesystem(self):
670 """Test without shared filesystem and no outputs."""
671 mock_workflow = unittest.mock.Mock()
672 mock_workflow.get_job_outputs.return_value = []
674 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
676 self.assertEqual(result, {"transfer_output_files": '""'})
678 def testWithAnOutputNoSharedFilesystem(self):
679 """Test without shared filesystem and single output file."""
680 mock_workflow = unittest.mock.Mock()
681 mock_workflow.get_job_outputs.return_value = [
682 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
683 ]
685 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
687 expected = {
688 "transfer_output_files": "output.txt",
689 "transfer_output_remaps": '"output.txt=/path/to/output.txt"',
690 }
691 self.assertEqual(result, expected)
693 def testWithOutputsNoSharedFilesystem(self):
694 """Test without shared filesystem and multiple output files."""
695 mock_workflow = unittest.mock.Mock()
696 mock_workflow.get_job_outputs.return_value = [
697 GenericWorkflowFile(name="output1.txt", src_uri="/path/output1.txt"),
698 GenericWorkflowFile(name="output2.txt", src_uri="/another/path/output2.txt"),
699 ]
701 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
703 expected = {
704 "transfer_output_files": "output1.txt,output2.txt",
705 "transfer_output_remaps": '"output1.txt=/path/output1.txt;output2.txt=/another/path/output2.txt"',
706 }
707 self.assertEqual(result, expected)
709 @unittest.mock.patch("lsst.ctrl.bps.htcondor.prepare_utils._LOG")
710 def testLogging(self, mock_log):
711 mock_workflow = unittest.mock.Mock()
712 mock_workflow.get_job_outputs.return_value = [
713 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt")
714 ]
716 prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix)
718 self.assertTrue(mock_log.debug.called)
719 debug_calls = mock_log.debug.call_args_list
720 self.assertTrue(any("src_uri=" in str(call) for call in debug_calls))
721 self.assertTrue(any("transfer_output_files=" in str(call) for call in debug_calls))
722 self.assertTrue(any("transfer_output_remaps=" in str(call) for call in debug_calls))
725class CreateJobTestCase(unittest.TestCase):
726 """Test _create_job function."""
728 def setUp(self):
729 self.generic_workflow = make_3_label_workflow("test1", True)
730 self.template = "{label}/{tract}/{patch}/{band}/{subfilter}/{physical_filter}/{visit}/{exposure}"
732 def testNoOverwrite(self):
733 cached_values = {
734 "bpsUseShared": True,
735 "overwriteJobFiles": False,
736 "memoryLimit": 491520,
737 "profile": {},
738 "attrs": {},
739 }
740 gwjob = self.generic_workflow.get_final()
741 out_prefix = "submit"
742 htc_job = prepare_utils._create_job(
743 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
744 )
745 self.assertEqual(htc_job.name, gwjob.name)
746 self.assertEqual(htc_job.label, gwjob.label)
747 self.assertIn("NumJobStarts", htc_job.cmds["output"])
748 self.assertIn("NumJobStarts", htc_job.cmds["error"])
749 self.assertNotIn("NumJobStarts", htc_job.cmds["log"])
750 self.assertTrue(htc_job.cmds["error"].endswith(".out"))
751 self.assertTrue(htc_job.cmds["output"].endswith(".out"))
752 self.assertTrue(htc_job.cmds["log"].endswith(".log"))
754 def testNodesetWithNoRequirements(self):
755 cached_values = {
756 "bpsUseShared": True,
757 "overwriteJobFiles": False,
758 "memoryLimit": 491520,
759 "profile": {},
760 "attrs": {},
761 "nodeset": "set1",
762 }
763 gwjob = self.generic_workflow.get_job("label1_10002_11")
764 out_prefix = "temp"
765 htc_job = prepare_utils._create_job(
766 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
767 )
768 self.assertEqual(htc_job.cmds["requirements"], '( Target.Nodeset == "set1" )')
769 self.assertEqual(htc_job.attrs["JobNodeset"], "set1")
771 def testNodesetWithRequirements(self):
772 cached_values = {
773 "bpsUseShared": True,
774 "overwriteJobFiles": False,
775 "memoryLimit": 491520,
776 "profile": {"requirements": "dummy_val == 3"},
777 "attrs": {},
778 "nodeset": "set1",
779 }
780 gwjob = self.generic_workflow.get_job("label1_10002_11")
781 out_prefix = "temp"
782 htc_job = prepare_utils._create_job(
783 self.template, cached_values, self.generic_workflow, gwjob, out_prefix
784 )
785 self.assertEqual(htc_job.cmds["requirements"], '(dummy_val == 3) && ( Target.Nodeset == "set1" )')
786 self.assertEqual(htc_job.attrs["JobNodeset"], "set1")
789class ReplaceWmsVarsTestCase(unittest.TestCase):
790 """Test _replace_wms_vars function."""
792 def testNoWmsVar(self):
793 orig_string = "whatever <Other:notThere> whatnot"
794 updated_string = prepare_utils._replace_wms_vars(orig_string)
795 self.assertEqual(orig_string, updated_string)
797 def testAttemptNum(self):
798 orig_string = "whatever <WMS:attemptNum> whatnot"
799 updated_string = prepare_utils._replace_wms_vars(orig_string)
800 self.assertEqual("whatever $$([NumJobStarts]) whatnot", updated_string)
802 def testUnrecognized(self):
803 orig_string = "whatever <WMS:notThere> whatnot"
804 with self.assertLogs(level="INFO") as cm_log:
805 with self.assertRaises(KeyError):
806 _ = prepare_utils._replace_wms_vars(orig_string)
807 self.assertRegex(cm_log.output[0], "Unrecognized WMS placeholder: notThere")
810class UpdateJobSummaryTestCase(unittest.TestCase):
811 """Test _update_job_summary function."""
813 def setUp(self):
814 self.run = "u_testuser_DM-53494_20260220T001651Z"
815 self.filename = f"{self.run}_ctrl.info.json"
816 self.data = {
817 "mycomputer": {
818 "24390.0": {
819 "ClusterId": 24390,
820 "GlobalJobId": "mycomputer#24390.0#1771546612",
821 "bps_run": f"{self.run}_ctrl",
822 "bps_isjob": "True",
823 "bps_payload": "DM-53494",
824 "bps_project": "dev",
825 "bps_runsite": "site1",
826 "bps_campaign": "ci_rc2",
827 "bps_operator": "testuser",
828 "bps_run_quanta": "",
829 "bps_job_summary": "buildQuantumGraph:1;preparePayloadWorkflow:1;dummyJob:1",
830 "bps_wms_service": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorService",
831 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_workflow.HTCondorWorkflow",
832 "bps_wms_config_path": "dagman.conf",
833 }
834 }
835 }
836 self.mapping = f"{self.run}:preparePayloadWorkflow"
837 self.add_summary = "pipetaskInit:1;isr:6;finalJob:1"
839 def testLazyMapping(self):
840 dag_info = deepcopy(self.data)
841 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = self.mapping
842 with temporaryDirectory() as tmp_dir:
843 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info)
845 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
847 _, results = lssthtc.read_dag_info(str(tmp_dir))
848 self.assertEqual(
849 results["mycomputer"]["24390.0"]["bps_job_summary"],
850 f"buildQuantumGraph:1;preparePayloadWorkflow:1;{self.add_summary};dummyJob:1",
851 )
853 def testNoLazyMapping(self):
854 # No bps_lazy_mapping at all
855 with temporaryDirectory() as tmp_dir:
856 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", self.data)
858 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir))
860 _, results = lssthtc.read_dag_info(str(tmp_dir))
861 self.assertEqual(
862 results["mycomputer"]["24390.0"]["bps_job_summary"],
863 f"{self.data['mycomputer']['24390.0']['bps_job_summary']};{self.add_summary}",
864 )
866 def testNoEntryLazyMapping(self):
867 # bps_lazy_mapping exists, but doesn't include this job
868 dag_info = deepcopy(self.data)
869 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = "other:preparePayloadWorkflow"
870 with temporaryDirectory() as tmp_dir:
871 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info)
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 )
882class ReplaceCmdVarsTestCase(unittest.TestCase):
883 """Test _replace_cmd_vars function."""
885 def testKeyError(self):
886 gwjob = GenericWorkflowJob("job1", "label1")
887 with self.assertLogs(level="DEBUG") as cm_log:
888 with self.assertRaisesRegex(KeyError, ".*notthere.*"):
889 _ = prepare_utils._replace_cmd_vars("{notthere}", gwjob)
890 self.assertRegex(cm_log.output[0], ".*replacement for 'notthere' not provided.*")
893class GenericWorkflowToHTCondorDAG(unittest.TestCase):
894 """Test _generic_workflow_to_htcondor_dag function."""
896 def testRegularWorkflow(self):
897 timestamp = "20260130T211713Z"
898 generic_workflow = make_3_label_workflow("test1", True)
899 config = BpsConfig(
900 {
901 "bpsUseShared": True,
902 "overwriteJobFiles": False,
903 "profile": {"requirements": "dummy_val == 3"},
904 "attrs": {},
905 "nodeset": "set1", # this shouldn't be used with auto-provisioning
906 "provisionResources": True,
907 "provisioning": {"provisioningMaxWallTime": 1200},
908 "bps_defined": {"timestamp": timestamp},
909 "saveHTCdot": True,
910 },
911 defaults=Config(HTC_DEFAULTS_URI),
912 )
914 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir")
915 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items())
916 self.assertIsNotNone(results.graph["final_job"])
917 self.assertTrue(is_isomorphic(results, generic_workflow))
918 self.assertTrue(results.graph["write_dot"])
920 def testLazyWorkflow(self):
921 timestamp = "20260130T211713Z"
922 generic_workflow = make_lazy_workflow("test1", True)
923 config = BpsConfig(
924 {
925 "bpsUseShared": True,
926 "overwriteJobFiles": False,
927 "profile": {"requirements": "dummy_val == 3"},
928 "attrs": {},
929 "nodeset": "set1", # this shouldn't be used with auto-provisioning
930 "provisionResources": True,
931 "provisioning": {"provisioningMaxWallTime": 1200},
932 "bps_defined": {"timestamp": timestamp},
933 },
934 defaults=Config(HTC_DEFAULTS_URI),
935 )
937 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir")
938 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items())
939 self.assertIsNotNone(results.graph["final_job"])
940 # Can't test isomorphic because HTCDag will have additional job for
941 # the lazy dagman job.
942 self.assertTrue(generic_workflow.nodes <= results.nodes)
943 self.assertFalse(results.graph["write_dot"])
946if __name__ == "__main__":
947 unittest.main()