Coverage for tests/test_prepare_utils.py: 100%

458 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-09 09:36 +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="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 def testPayloadCommandEnvironmentShell(self): 

283 # Exports in commands, no environment in jobcmds 

284 gw_exec = GenericWorkflowExec("test_exec", "<ENV:CTRL_BPS_DIR>/bin/pipetask") 

285 gw, gwjob = self._make_job(executable=gw_exec, arguments="go") 

286 gwjob.environment = {"TEST_ENV_VAR": "<ENV:CTRL_BPS_DIR>/tests"} 

287 cached_vals = { 

288 "bpsUseShared": True, 

289 "bpsMakeCommand": False, 

290 "bpsUseHTCEnvironment": False, 

291 "payloadCommand": "{gwjobExports} {gwjobCommand}", 

292 } 

293 jobcmds = prepare_utils._translate_command_line(cached_vals, gw, gwjob) 

294 self.assertIn("export TEST_ENV_VAR='${CTRL_BPS_DIR}/tests';", jobcmds["arguments"]) 

295 self.assertNotIn("environment", jobcmds) 

296 

297 

298class TranslateDagCmdsTestCase(unittest.TestCase): 

299 """Test _translate_dag_cmds method.""" 

300 

301 def setUp(self): 

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

303 

304 def testPriority(self): 

305 gwjob = GenericWorkflowJob("priority", "label1", executable=self.gw_exec) 

306 gwjob.priority = 100 

307 dag_commands = prepare_utils._translate_dag_cmds(gwjob) 

308 self.assertEqual(dag_commands["priority"], 100) 

309 

310 

311class GroupToSubdagTestCase(unittest.TestCase): 

312 """Test _group_to_subdag function.""" 

313 

314 def testBlocking(self): 

315 gw = make_3_label_workflow_groups_sort("test1", True) 

316 gwjob = gw.get_job("group_order1_10001") 

317 config = BpsConfig( 

318 {}, 

319 search_order=BPS_SEARCH_ORDER, 

320 defaults=BPS_DEFAULTS, 

321 ) 

322 

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

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

325 

326 

327class GatherSiteValuesTestCase(unittest.TestCase): 

328 """Test _gather_site_values function.""" 

329 

330 def testAllThere(self): 

331 config = BpsConfig( 

332 {}, 

333 search_order=BPS_SEARCH_ORDER, 

334 defaults=BPS_DEFAULTS, 

335 ) 

336 compute_site = "notThere" 

337 results = prepare_utils._gather_site_values(config, compute_site) 

338 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"]) 

339 

340 def testNotSpecified(self): 

341 config = BpsConfig( 

342 {}, 

343 search_order=BPS_SEARCH_ORDER, 

344 defaults=BPS_DEFAULTS, 

345 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

346 ) 

347 compute_site = "notThere" 

348 results = prepare_utils._gather_site_values(config, compute_site) 

349 self.assertEqual(results["memoryLimit"], BPS_DEFAULTS["memoryLimit"]) 

350 

351 def testAttrsProfile(self): 

352 test_values = { 

353 "bpsNodeset": "DEVSET", 

354 "site": { 

355 "mycomputer": { 

356 "profile": { 

357 "condor": { 

358 "requirements": '( TARGET.Nodeset == "{bpsNodeset}" )', 

359 "+JobNodeset": "{bpsNodeset}", 

360 } 

361 } 

362 } 

363 }, 

364 } 

365 config = BpsConfig( 

366 test_values, 

367 search_order=BPS_SEARCH_ORDER, 

368 defaults=BPS_DEFAULTS, 

369 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

370 ) 

371 results = prepare_utils._gather_site_values(config, "mycomputer") 

372 self.assertEqual(results["profile"], {"requirements": '( TARGET.Nodeset == "DEVSET" )'}) 

373 self.assertEqual(results["attrs"], {"JobNodeset": "DEVSET"}) 

374 

375 

376class GatherLabelValuesTestCase(unittest.TestCase): 

377 """Test _gather_labels_values function.""" 

378 

379 def testClusterLabel(self): 

380 # Test cluster value overrides pipetask. 

381 config = BpsConfig( 

382 { 

383 "cluster": { 

384 "label1": { 

385 "releaseExpr": "cluster_val", 

386 "overwriteJobFiles": False, 

387 "profile": {"condor": {"prof_val1": 3}}, 

388 } 

389 }, 

390 "pipetask": {"label1": {"releaseExpr": "pipetask_val"}}, 

391 "site": {"site1": {}}, 

392 }, 

393 search_order=BPS_SEARCH_ORDER, 

394 defaults=BPS_DEFAULTS, 

395 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

396 ) 

397 results = prepare_utils._gather_label_values(config, "label1") 

398 self.assertEqual( 

399 results, 

400 { 

401 "attrs": {}, 

402 "profile": {"prof_val1": 3}, 

403 "releaseExpr": "cluster_val", 

404 "overwriteJobFiles": False, 

405 "bpsMakeCommand": True, 

406 "bpsUseHTCEnvironment": True, 

407 "bpsUseShared": True, 

408 "memoryLimit": 491520, 

409 }, 

410 ) 

411 

412 def testPipetaskLabel(self): 

413 label = "label1" 

414 config = BpsConfig( 

415 { 

416 "pipetask": { 

417 "label1": { 

418 "releaseExpr": "pipetask_val", 

419 "overwriteJobFiles": False, 

420 "profile": {"condor": {"prof_val1": 3}}, 

421 } 

422 }, 

423 "site": {"site1": {}}, 

424 }, 

425 search_order=BPS_SEARCH_ORDER, 

426 defaults=BPS_DEFAULTS, 

427 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

428 ) 

429 results = prepare_utils._gather_label_values(config, label) 

430 self.assertEqual( 

431 results, 

432 { 

433 "attrs": {}, 

434 "bpsMakeCommand": True, 

435 "bpsUseHTCEnvironment": True, 

436 "bpsUseShared": True, 

437 "memoryLimit": 491520, 

438 "overwriteJobFiles": False, 

439 "profile": {"prof_val1": 3}, 

440 "releaseExpr": "pipetask_val", 

441 }, 

442 ) 

443 

444 def testNoSection(self): 

445 label = "notThere" 

446 config = BpsConfig( 

447 {"site": {"site1": {}}}, 

448 search_order=BPS_SEARCH_ORDER, 

449 defaults=BPS_DEFAULTS, 

450 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

451 ) 

452 results = prepare_utils._gather_label_values(config, label) 

453 self.assertEqual( 

454 results, 

455 { 

456 "attrs": {}, 

457 "profile": {}, 

458 "overwriteJobFiles": True, 

459 "bpsMakeCommand": True, 

460 "bpsUseHTCEnvironment": True, 

461 "bpsUseShared": True, 

462 "memoryLimit": 491520, 

463 }, 

464 ) 

465 

466 def testNoOverwriteSpecified(self): 

467 label = "notthere" 

468 config = BpsConfig( 

469 {"site": {"site1": {}}, "memoryLimit": 491520}, 

470 search_order=BPS_SEARCH_ORDER, 

471 defaults={}, 

472 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

473 ) 

474 results = prepare_utils._gather_label_values(config, label) 

475 self.assertEqual( 

476 results, 

477 { 

478 "attrs": {}, 

479 "profile": {}, 

480 "overwriteJobFiles": True, 

481 "bpsMakeCommand": True, 

482 "bpsUseHTCEnvironment": True, 

483 "bpsUseShared": False, 

484 "memoryLimit": 491520, 

485 }, 

486 ) 

487 

488 def testFinalJob(self): 

489 label = "finalJob" 

490 config = BpsConfig( 

491 {"site": {"site1": {}}, "finalJob": {"profile": {"condor": {"prof_val2": 6, "+attr_val1": 5}}}}, 

492 search_order=BPS_SEARCH_ORDER, 

493 defaults=BPS_DEFAULTS, 

494 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

495 ) 

496 results = prepare_utils._gather_label_values(config, label) 

497 self.assertEqual( 

498 results, 

499 { 

500 "attrs": {"attr_val1": 5}, 

501 "profile": {"prof_val2": 6}, 

502 "overwriteJobFiles": False, 

503 "bpsMakeCommand": True, 

504 "bpsUseHTCEnvironment": True, 

505 "bpsUseShared": True, 

506 "memoryLimit": 491520, 

507 }, 

508 ) 

509 

510 def testGlobalNodeset(self): 

511 config = BpsConfig( 

512 {"nodeset": "global_node_set_{campaign}", "campaign": "DRP"}, 

513 search_order=BPS_SEARCH_ORDER, 

514 defaults=BPS_DEFAULTS, 

515 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

516 ) 

517 results = prepare_utils._gather_label_values(config, "label1") 

518 self.assertEqual(results["nodeset"], "global_node_set_DRP") 

519 

520 def testSiteNodeset(self): 

521 config = BpsConfig( 

522 { 

523 "nodeset": "global_node_set_{campaign}", 

524 "campaign": "DRP", 

525 "site": {"fr": {"nodeset": "fr_node_set_{campaign}", "siteVar": "frSiteVal"}}, 

526 "computeSite": "fr", 

527 }, 

528 search_order=BPS_SEARCH_ORDER, 

529 defaults=BPS_DEFAULTS, 

530 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

531 ) 

532 results = prepare_utils._gather_label_values(config, "label1") 

533 self.assertEqual(results["nodeset"], "fr_node_set_DRP") 

534 self.assertEqual(results["siteVar"], "frSiteVal") 

535 

536 def testBpsMakeCommandFalse(self): 

537 config = BpsConfig( 

538 { 

539 "bpsMakeCommand": False, 

540 }, 

541 search_order=BPS_SEARCH_ORDER, 

542 defaults=BPS_DEFAULTS, 

543 wms_service_class_fqn="lsst.ctrl.bps.htcondor.HTCondorService", 

544 ) 

545 results = prepare_utils._gather_label_values(config, "label1") 

546 self.assertIn("payloadCommand", results) 

547 self.assertIn("gwjobCommand", results["payloadCommand"]) 

548 

549 

550class CreateCheckJobTestCase(unittest.TestCase): 

551 """Test _create_check_job function.""" 

552 

553 def testSuccess(self): 

554 group_job_name = "group_order1_val1a" 

555 job_label = "order1" 

556 job = prepare_utils._create_check_job(group_job_name, job_label, {}) 

557 self.assertIn(group_job_name, job.name) 

558 self.assertEqual(job.label, job_label) 

559 self.assertIn("check_group_status.sub", job.subfile) 

560 self.assertNotIn("job_nodeset", job.dagcmds["vars"]) 

561 

562 def testNodeSetSuccess(self): 

563 group_job_name = "group_order1_val1a" 

564 job_label = "order1" 

565 job = prepare_utils._create_check_job(group_job_name, job_label, {"nodeset": "custom_nodeset"}) 

566 self.assertIn(group_job_name, job.name) 

567 self.assertEqual(job.label, job_label) 

568 self.assertIn("check_group_status.sub", job.subfile) 

569 self.assertIn("job_nodeset", job.dagcmds["vars"]) 

570 

571 

572class CreatePeriodicReleaseExprTestCase(unittest.TestCase): 

573 """Test _create_periodic_release_expr function.""" 

574 

575 def setUp(self): 

576 self.maxDiff = None 

577 

578 def testNoReleaseExpr(self): 

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

580 self.assertEqual(results, "") 

581 

582 def testMultiplierNone(self): 

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

584 self.assertEqual(results, "") 

585 

586 def testJustMemoryReleaseExpr(self): 

587 self.maxDiff = None # so test error shows entire strings 

588 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "") 

589 truth = ( 

590 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && " 

591 "(HoldReasonCode =?= 12 || " 

592 "(HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || " 

593 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && " 

594 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768)" 

595 ) 

596 self.assertEqual(results, truth) 

597 

598 def testJustUserReleaseExpr(self): 

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

600 truth = ( 

601 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && " 

602 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)" 

603 ) 

604 self.assertEqual(results, truth) 

605 

606 def testJustUserReleaseExprMultiplierNone(self): 

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

608 truth = ( 

609 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && " 

610 "(HoldReasonCode =?= 12 || HoldReasonCode =!= 1 && True)" 

611 ) 

612 self.assertEqual(results, truth) 

613 

614 def testMemoryAndUserReleaseExpr(self): 

615 self.maxDiff = None # so test error shows entire strings 

616 results = prepare_utils._create_periodic_release_expr(2048, 2, 32768, "True") 

617 truth = ( 

618 "JobStatus == 5 && NumJobStarts <= JobMaxRetries && " 

619 "(HoldReasonCode =?= 12 || (HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || " 

620 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && " 

621 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) < 32768 || " 

622 "HoldReasonCode =!= 1 && True)" 

623 ) 

624 self.assertEqual(results, truth) 

625 

626 

627class CreatePeriodicRemoveExprTestCase(unittest.TestCase): 

628 """Test _create_periodic_release_expr function.""" 

629 

630 def testBasicRemoveExpr(self): 

631 """Function assumes only called if max_retries >= 0.""" 

632 results = prepare_utils._create_periodic_remove_expr(2048, 1, 32768) 

633 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)" 

634 self.assertEqual(results, truth) 

635 

636 def testBasicRemoveExprMultiplierNone(self): 

637 """Function assumes only called if max_retries >= 0.""" 

638 results = prepare_utils._create_periodic_remove_expr(2048, None, 32768) 

639 truth = "JobStatus == 5 && (NumJobStarts > JobMaxRetries)" 

640 self.assertEqual(results, truth) 

641 

642 def testMemoryRemoveExpr(self): 

643 self.maxDiff = None # so test error shows entire strings 

644 results = prepare_utils._create_periodic_remove_expr(2048, 2, 32768) 

645 truth = ( 

646 "JobStatus == 5 && (NumJobStarts > JobMaxRetries || " 

647 "((HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 || " 

648 "HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && " 

649 "min({int(2048 * pow(2, NumJobStarts - 1)), 32768}) == 32768))" 

650 ) 

651 self.assertEqual(results, truth) 

652 

653 

654class HandleJobOutputsTestCase(unittest.TestCase): 

655 """Test _handle_job_outputs function.""" 

656 

657 def setUp(self): 

658 self.job_name = "test_job" 

659 self.out_prefix = "/test/prefix" 

660 

661 def tearDown(self): 

662 pass 

663 

664 def testNoOutputsSharedFilesystem(self): 

665 """Test with shared filesystem and no outputs.""" 

666 mock_workflow = unittest.mock.Mock() 

667 mock_workflow.get_job_outputs.return_value = [] 

668 

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

670 

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

672 

673 def testWithOutputsSharedFilesystem(self): 

674 """Test with shared filesystem and outputs present (still empty).""" 

675 mock_workflow = unittest.mock.Mock() 

676 mock_workflow.get_job_outputs.return_value = [ 

677 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt") 

678 ] 

679 

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

681 

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

683 

684 def testNoOutputsNoSharedFilesystem(self): 

685 """Test without shared filesystem and no outputs.""" 

686 mock_workflow = unittest.mock.Mock() 

687 mock_workflow.get_job_outputs.return_value = [] 

688 

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

690 

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

692 

693 def testWithAnOutputNoSharedFilesystem(self): 

694 """Test without shared filesystem and single output file.""" 

695 mock_workflow = unittest.mock.Mock() 

696 mock_workflow.get_job_outputs.return_value = [ 

697 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt") 

698 ] 

699 

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

701 

702 expected = { 

703 "transfer_output_files": "output.txt", 

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

705 } 

706 self.assertEqual(result, expected) 

707 

708 def testWithOutputsNoSharedFilesystem(self): 

709 """Test without shared filesystem and multiple output files.""" 

710 mock_workflow = unittest.mock.Mock() 

711 mock_workflow.get_job_outputs.return_value = [ 

712 GenericWorkflowFile(name="output1.txt", src_uri="/path/output1.txt"), 

713 GenericWorkflowFile(name="output2.txt", src_uri="/another/path/output2.txt"), 

714 ] 

715 

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

717 

718 expected = { 

719 "transfer_output_files": "output1.txt,output2.txt", 

720 "transfer_output_remaps": '"output1.txt=/path/output1.txt;output2.txt=/another/path/output2.txt"', 

721 } 

722 self.assertEqual(result, expected) 

723 

724 @unittest.mock.patch("lsst.ctrl.bps.htcondor.prepare_utils._LOG") 

725 def testLogging(self, mock_log): 

726 mock_workflow = unittest.mock.Mock() 

727 mock_workflow.get_job_outputs.return_value = [ 

728 GenericWorkflowFile(name="output.txt", src_uri="/path/to/output.txt") 

729 ] 

730 

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

732 

733 self.assertTrue(mock_log.debug.called) 

734 debug_calls = mock_log.debug.call_args_list 

735 self.assertTrue(any("src_uri=" in str(call) for call in debug_calls)) 

736 self.assertTrue(any("transfer_output_files=" in str(call) for call in debug_calls)) 

737 self.assertTrue(any("transfer_output_remaps=" in str(call) for call in debug_calls)) 

738 

739 

740class CreateJobTestCase(unittest.TestCase): 

741 """Test _create_job function.""" 

742 

743 def setUp(self): 

744 self.generic_workflow = make_3_label_workflow("test1", True) 

745 self.template = "{label}/{tract}/{patch}/{band}/{subfilter}/{physical_filter}/{visit}/{exposure}" 

746 

747 def testNoOverwrite(self): 

748 cached_values = { 

749 "bpsUseShared": True, 

750 "overwriteJobFiles": False, 

751 "memoryLimit": 491520, 

752 "profile": {}, 

753 "attrs": {}, 

754 } 

755 gwjob = self.generic_workflow.get_final() 

756 out_prefix = "submit" 

757 htc_job = prepare_utils._create_job( 

758 self.template, cached_values, self.generic_workflow, gwjob, out_prefix 

759 ) 

760 self.assertEqual(htc_job.name, gwjob.name) 

761 self.assertEqual(htc_job.label, gwjob.label) 

762 self.assertIn("NumJobStarts", htc_job.cmds["output"]) 

763 self.assertIn("NumJobStarts", htc_job.cmds["error"]) 

764 self.assertNotIn("NumJobStarts", htc_job.cmds["log"]) 

765 self.assertTrue(htc_job.cmds["error"].endswith(".out")) 

766 self.assertTrue(htc_job.cmds["output"].endswith(".out")) 

767 self.assertTrue(htc_job.cmds["log"].endswith(".log")) 

768 

769 def testNodesetWithNoRequirements(self): 

770 cached_values = { 

771 "bpsUseShared": True, 

772 "overwriteJobFiles": False, 

773 "memoryLimit": 491520, 

774 "profile": {}, 

775 "attrs": {}, 

776 "nodeset": "set1", 

777 } 

778 gwjob = self.generic_workflow.get_job("label1_10002_11") 

779 out_prefix = "temp" 

780 htc_job = prepare_utils._create_job( 

781 self.template, cached_values, self.generic_workflow, gwjob, out_prefix 

782 ) 

783 self.assertEqual(htc_job.cmds["requirements"], '( Target.Nodeset == "set1" )') 

784 self.assertEqual(htc_job.attrs["JobNodeset"], "set1") 

785 

786 def testNodesetWithRequirements(self): 

787 cached_values = { 

788 "bpsUseShared": True, 

789 "overwriteJobFiles": False, 

790 "memoryLimit": 491520, 

791 "profile": {"requirements": "dummy_val == 3"}, 

792 "attrs": {}, 

793 "nodeset": "set1", 

794 } 

795 gwjob = self.generic_workflow.get_job("label1_10002_11") 

796 out_prefix = "temp" 

797 htc_job = prepare_utils._create_job( 

798 self.template, cached_values, self.generic_workflow, gwjob, out_prefix 

799 ) 

800 self.assertEqual(htc_job.cmds["requirements"], '(dummy_val == 3) && ( Target.Nodeset == "set1" )') 

801 self.assertEqual(htc_job.attrs["JobNodeset"], "set1") 

802 

803 

804class ReplaceWmsVarsTestCase(unittest.TestCase): 

805 """Test _replace_wms_vars function.""" 

806 

807 def testNoWmsVar(self): 

808 orig_string = "whatever <Other:notThere> whatnot" 

809 updated_string = prepare_utils._replace_wms_vars(orig_string) 

810 self.assertEqual(orig_string, updated_string) 

811 

812 def testAttemptNum(self): 

813 orig_string = "whatever <WMS:attemptNum> whatnot" 

814 updated_string = prepare_utils._replace_wms_vars(orig_string) 

815 self.assertEqual("whatever $$([NumJobStarts]) whatnot", updated_string) 

816 

817 def testUnrecognized(self): 

818 orig_string = "whatever <WMS:notThere> whatnot" 

819 with self.assertLogs(level="INFO") as cm_log: 

820 with self.assertRaises(KeyError): 

821 _ = prepare_utils._replace_wms_vars(orig_string) 

822 self.assertRegex(cm_log.output[0], "Unrecognized WMS placeholder: notThere") 

823 

824 

825class UpdateJobSummaryTestCase(unittest.TestCase): 

826 """Test _update_job_summary function.""" 

827 

828 def setUp(self): 

829 self.run = "u_testuser_DM-53494_20260220T001651Z" 

830 self.filename = f"{self.run}_ctrl.info.json" 

831 self.data = { 

832 "mycomputer": { 

833 "24390.0": { 

834 "ClusterId": 24390, 

835 "GlobalJobId": "mycomputer#24390.0#1771546612", 

836 "bps_run": f"{self.run}_ctrl", 

837 "bps_isjob": "True", 

838 "bps_payload": "DM-53494", 

839 "bps_project": "dev", 

840 "bps_runsite": "site1", 

841 "bps_campaign": "ci_rc2", 

842 "bps_operator": "testuser", 

843 "bps_run_quanta": "", 

844 "bps_job_summary": "buildQuantumGraph:1;preparePayloadWorkflow:1;dummyJob:1", 

845 "bps_wms_service": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorService", 

846 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_workflow.HTCondorWorkflow", 

847 "bps_wms_config_path": "dagman.conf", 

848 } 

849 } 

850 } 

851 self.mapping = f"{self.run}:preparePayloadWorkflow" 

852 self.add_summary = "pipetaskInit:1;isr:6;finalJob:1" 

853 

854 def testLazyMapping(self): 

855 dag_info = deepcopy(self.data) 

856 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = self.mapping 

857 with temporaryDirectory() as tmp_dir: 

858 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info) 

859 

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

861 

862 _, results = lssthtc.read_dag_info(str(tmp_dir)) 

863 self.assertEqual( 

864 results["mycomputer"]["24390.0"]["bps_job_summary"], 

865 f"buildQuantumGraph:1;preparePayloadWorkflow:1;{self.add_summary};dummyJob:1", 

866 ) 

867 

868 def testNoLazyMapping(self): 

869 # No bps_lazy_mapping at all 

870 with temporaryDirectory() as tmp_dir: 

871 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", self.data) 

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 def testNoEntryLazyMapping(self): 

882 # bps_lazy_mapping exists, but doesn't include this job 

883 dag_info = deepcopy(self.data) 

884 dag_info["mycomputer"]["24390.0"]["bps_lazy_mapping"] = "other:preparePayloadWorkflow" 

885 with temporaryDirectory() as tmp_dir: 

886 lssthtc.write_dag_info(f"{tmp_dir}/{self.filename}", dag_info) 

887 

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

889 

890 _, results = lssthtc.read_dag_info(str(tmp_dir)) 

891 self.assertEqual( 

892 results["mycomputer"]["24390.0"]["bps_job_summary"], 

893 f"{self.data['mycomputer']['24390.0']['bps_job_summary']};{self.add_summary}", 

894 ) 

895 

896 

897class ReplaceCmdVarsTestCase(unittest.TestCase): 

898 """Test _replace_cmd_vars function.""" 

899 

900 def testKeyError(self): 

901 gwjob = GenericWorkflowJob("job1", "label1") 

902 with self.assertLogs(level="DEBUG") as cm_log: 

903 with self.assertRaisesRegex(KeyError, ".*notthere.*"): 

904 _ = prepare_utils._replace_cmd_vars("{notthere}", gwjob) 

905 self.assertRegex(cm_log.output[0], ".*replacement for 'notthere' not provided.*") 

906 

907 

908class GenericWorkflowToHTCondorDAG(unittest.TestCase): 

909 """Test _generic_workflow_to_htcondor_dag function.""" 

910 

911 def testRegularWorkflow(self): 

912 timestamp = "20260130T211713Z" 

913 generic_workflow = make_3_label_workflow("test1", True) 

914 config = BpsConfig( 

915 { 

916 "bpsUseShared": True, 

917 "overwriteJobFiles": False, 

918 "profile": {"requirements": "dummy_val == 3"}, 

919 "attrs": {}, 

920 "nodeset": "set1", # this shouldn't be used with auto-provisioning 

921 "provisionResources": True, 

922 "provisioning": {"provisioningMaxWallTime": 1200}, 

923 "bps_defined": {"timestamp": timestamp}, 

924 "saveHTCdot": True, 

925 }, 

926 defaults=Config(HTC_DEFAULTS_URI), 

927 ) 

928 

929 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir") 

930 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items()) 

931 self.assertIsNotNone(results.graph["final_job"]) 

932 self.assertTrue(is_isomorphic(results, generic_workflow)) 

933 self.assertTrue(results.graph["write_dot"]) 

934 

935 def testLazyWorkflow(self): 

936 timestamp = "20260130T211713Z" 

937 generic_workflow = make_lazy_workflow("test1", True) 

938 config = BpsConfig( 

939 { 

940 "bpsUseShared": True, 

941 "overwriteJobFiles": False, 

942 "profile": {"requirements": "dummy_val == 3"}, 

943 "attrs": {}, 

944 "nodeset": "set1", # this shouldn't be used with auto-provisioning 

945 "provisionResources": True, 

946 "provisioning": {"provisioningMaxWallTime": 1200}, 

947 "bps_defined": {"timestamp": timestamp}, 

948 }, 

949 defaults=Config(HTC_DEFAULTS_URI), 

950 ) 

951 

952 results = prepare_utils._generic_workflow_to_htcondor_dag(config, generic_workflow, "/mock_dir") 

953 self.assertTrue(generic_workflow.run_attrs.items() <= results.graph["attr"].items()) 

954 self.assertIsNotNone(results.graph["final_job"]) 

955 # Can't test isomorphic because HTCDag will have additional job for 

956 # the lazy dagman job. 

957 self.assertTrue(generic_workflow.nodes <= results.nodes) 

958 self.assertFalse(results.graph["write_dot"]) 

959 

960 

961if __name__ == "__main__": 

962 unittest.main()