Coverage for tests/test_prepare_utils.py: 100%

466 statements  

« prev     ^ index     » next       coverage.py v7.16.2, created at 2026-09-29 09:37 +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/>. 

27 

28"""Unit tests for prepare utility functions.""" 

29 

30import logging 

31import os 

32import unittest 

33from copy import deepcopy 

34 

35from networkx import is_isomorphic 

36 

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 

55 

56logger = logging.getLogger("lsst.ctrl.bps.htcondor") 

57 

58TESTDIR = os.path.abspath(os.path.dirname(__file__)) 

59 

60 

61class TranslateJobCmdsTestCase(unittest.TestCase): 

62 """Test _translate_job_cmds method.""" 

63 

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 } 

73 

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) 

79 

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) 

85 

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})") 

91 

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)) 

98 

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'") 

104 

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'") 

110 

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'") 

116 

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'""") 

122 

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'") 

128 

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) 

142 

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) 

152 

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) 

162 

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) 

173 

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) 

184 

185 

186class TranslateCommandLineTestCase(unittest.TestCase): 

187 """Test _translate_command_line method.""" 

188 

189 def setUp(self): 

190 self.gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask") 

191 self.cached_vals = {"bpsUseShared": True, "bpsMakeCommand": True, "bpsUseHTCEnvironment": True} 

192 

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 

200 

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) 

209 

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") 

216 

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") 

224 

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") 

233 

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'") 

247 

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 '") 

261 

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'") 

271 

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'") 

279 

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) 

289 

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'") 

295 

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) 

310 

311 

312class TranslateDagCmdsTestCase(unittest.TestCase): 

313 """Test _translate_dag_cmds method.""" 

314 

315 def setUp(self): 

316 self.gw_exec = GenericWorkflowExec("test_exec", "/dummy/dir/pipetask") 

317 

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) 

323 

324 

325class GroupToSubdagTestCase(unittest.TestCase): 

326 """Test _group_to_subdag function.""" 

327 

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 ) 

336 

337 htc_job = prepare_utils._group_to_subdag(config, gwjob, "the_prefix") 

338 self.assertEqual(len(htc_job.subdag), len(gwjob)) 

339 

340 

341class GatherSiteValuesTestCase(unittest.TestCase): 

342 """Test _gather_site_values function.""" 

343 

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"]) 

353 

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"]) 

364 

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"}) 

388 

389 

390class GatherLabelValuesTestCase(unittest.TestCase): 

391 """Test _gather_labels_values function.""" 

392 

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 ) 

425 

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 ) 

457 

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 ) 

479 

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 ) 

501 

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 ) 

523 

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") 

533 

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") 

549 

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"]) 

562 

563 

564class CreateCheckJobTestCase(unittest.TestCase): 

565 """Test _create_check_job function.""" 

566 

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"]) 

575 

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"]) 

584 

585 

586class CreatePeriodicReleaseExprTestCase(unittest.TestCase): 

587 """Test _create_periodic_release_expr function.""" 

588 

589 def setUp(self): 

590 self.maxDiff = None 

591 

592 def testNoReleaseExpr(self): 

593 results = prepare_utils._create_periodic_release_expr(2048, 1, 32768, "") 

594 self.assertEqual(results, "") 

595 

596 def testMultiplierNone(self): 

597 results = prepare_utils._create_periodic_release_expr(2048, None, 32768, "") 

598 self.assertEqual(results, "") 

599 

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) 

611 

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) 

619 

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) 

627 

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) 

639 

640 

641class CreatePeriodicRemoveExprTestCase(unittest.TestCase): 

642 """Test _create_periodic_release_expr function.""" 

643 

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) 

649 

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) 

655 

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) 

666 

667 

668class HandleJobOutputsTestCase(unittest.TestCase): 

669 """Test _handle_job_outputs function.""" 

670 

671 def setUp(self): 

672 self.job_name = "test_job" 

673 self.out_prefix = "/test/prefix" 

674 

675 def tearDown(self): 

676 pass 

677 

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 = [] 

682 

683 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix) 

684 

685 self.assertEqual(result, {"transfer_output_files": '""'}) 

686 

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 ] 

693 

694 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, True, self.out_prefix) 

695 

696 self.assertEqual(result, {"transfer_output_files": '""'}) 

697 

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 = [] 

702 

703 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix) 

704 

705 self.assertEqual(result, {"transfer_output_files": '""'}) 

706 

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 ] 

713 

714 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix) 

715 

716 expected = { 

717 "transfer_output_files": "output.txt", 

718 "transfer_output_remaps": '"output.txt=/path/to/output.txt"', 

719 } 

720 self.assertEqual(result, expected) 

721 

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 ] 

729 

730 result = prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix) 

731 

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) 

737 

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 ] 

744 

745 prepare_utils._handle_job_outputs(mock_workflow, self.job_name, False, self.out_prefix) 

746 

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)) 

752 

753 

754class CreateJobTestCase(unittest.TestCase): 

755 """Test _create_job function.""" 

756 

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}" 

760 

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")) 

782 

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") 

799 

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") 

816 

817 

818class ReplaceWmsVarsTestCase(unittest.TestCase): 

819 """Test _replace_wms_vars function.""" 

820 

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) 

825 

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) 

830 

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") 

837 

838 

839class UpdateJobSummaryTestCase(unittest.TestCase): 

840 """Test _update_job_summary function.""" 

841 

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" 

867 

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) 

873 

874 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir)) 

875 

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 ) 

881 

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) 

886 

887 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir)) 

888 

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 ) 

894 

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) 

901 

902 prepare_utils._update_job_summary(self.run, self.add_summary, str(tmp_dir)) 

903 

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 ) 

909 

910 

911class ReplaceCmdVarsTestCase(unittest.TestCase): 

912 """Test _replace_cmd_vars function.""" 

913 

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.*") 

920 

921 

922class GenericWorkflowToHTCondorDAG(unittest.TestCase): 

923 """Test _generic_workflow_to_htcondor_dag function.""" 

924 

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 ) 

942 

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"]) 

948 

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 ) 

965 

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"]) 

973 

974 

975if __name__ == "__main__": 

976 unittest.main()