Coverage for tests/test_prepare_utils.py: 100%

450 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-01 09:26 +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='${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="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'") 

247 

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

257 

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

265 

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) 

275 

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

281 

282 

283class TranslateDagCmdsTestCase(unittest.TestCase): 

284 """Test _translate_dag_cmds method.""" 

285 

286 def setUp(self): 

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

288 

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) 

294 

295 

296class GroupToSubdagTestCase(unittest.TestCase): 

297 """Test _group_to_subdag function.""" 

298 

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 ) 

307 

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

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

310 

311 

312class GatherSiteValuesTestCase(unittest.TestCase): 

313 """Test _gather_site_values function.""" 

314 

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

324 

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

335 

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

359 

360 

361class GatherLabelValuesTestCase(unittest.TestCase): 

362 """Test _gather_labels_values function.""" 

363 

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 ) 

396 

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 ) 

428 

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 ) 

450 

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 ) 

472 

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 ) 

494 

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

504 

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

520 

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

533 

534 

535class CreateCheckJobTestCase(unittest.TestCase): 

536 """Test _create_check_job function.""" 

537 

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

546 

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

555 

556 

557class CreatePeriodicReleaseExprTestCase(unittest.TestCase): 

558 """Test _create_periodic_release_expr function.""" 

559 

560 def setUp(self): 

561 self.maxDiff = None 

562 

563 def testNoReleaseExpr(self): 

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

565 self.assertEqual(results, "") 

566 

567 def testMultiplierNone(self): 

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

569 self.assertEqual(results, "") 

570 

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) 

582 

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) 

590 

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) 

598 

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) 

610 

611 

612class CreatePeriodicRemoveExprTestCase(unittest.TestCase): 

613 """Test _create_periodic_release_expr function.""" 

614 

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) 

620 

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) 

626 

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) 

637 

638 

639class HandleJobOutputsTestCase(unittest.TestCase): 

640 """Test _handle_job_outputs function.""" 

641 

642 def setUp(self): 

643 self.job_name = "test_job" 

644 self.out_prefix = "/test/prefix" 

645 

646 def tearDown(self): 

647 pass 

648 

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

653 

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

655 

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

657 

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 ] 

664 

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

666 

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

668 

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

673 

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

675 

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

677 

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 ] 

684 

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

686 

687 expected = { 

688 "transfer_output_files": "output.txt", 

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

690 } 

691 self.assertEqual(result, expected) 

692 

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 ] 

700 

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

702 

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) 

708 

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 ] 

715 

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

717 

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

723 

724 

725class CreateJobTestCase(unittest.TestCase): 

726 """Test _create_job function.""" 

727 

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

731 

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

753 

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

770 

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

787 

788 

789class ReplaceWmsVarsTestCase(unittest.TestCase): 

790 """Test _replace_wms_vars function.""" 

791 

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) 

796 

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) 

801 

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

808 

809 

810class UpdateJobSummaryTestCase(unittest.TestCase): 

811 """Test _update_job_summary function.""" 

812 

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" 

838 

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) 

844 

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

846 

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 ) 

852 

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) 

857 

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

859 

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 ) 

865 

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) 

872 

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

874 

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 ) 

880 

881 

882class ReplaceCmdVarsTestCase(unittest.TestCase): 

883 """Test _replace_cmd_vars function.""" 

884 

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

891 

892 

893class GenericWorkflowToHTCondorDAG(unittest.TestCase): 

894 """Test _generic_workflow_to_htcondor_dag function.""" 

895 

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 ) 

913 

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

919 

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 ) 

936 

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

944 

945 

946if __name__ == "__main__": 

947 unittest.main()