Coverage for tests/test_lssthtc.py: 99%

829 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-14 09:46 +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"""Unit tests for classes and functions in lssthtc.py.""" 

28 

29import io 

30import logging 

31import os 

32import pathlib 

33import stat 

34import tempfile 

35import unittest 

36from shutil import copy2, copytree, ignore_patterns, rmtree, which 

37 

38import htcondor 

39from dag_test_utils import make_lazy_dag 

40 

41from lsst.ctrl.bps import BpsConfig 

42from lsst.ctrl.bps.bps_utils import chdir 

43from lsst.ctrl.bps.htcondor import dagman_configurator, htcondor_config, lssthtc 

44from lsst.daf.butler import Config 

45from lsst.utils.tests import temporaryDirectory 

46 

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

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

49 

50 

51class TestLsstHtc(unittest.TestCase): 

52 """Test basic usage.""" 

53 

54 def testHtcEscapeInt(self): 

55 self.assertEqual(lssthtc.htc_escape(100), 100) 

56 

57 def testHtcEscapeDouble(self): 

58 self.assertEqual(lssthtc.htc_escape('"double"'), '""double""') 

59 

60 def testHtcEscapeSingle(self): 

61 self.assertEqual(lssthtc.htc_escape("'single'"), "''single''") 

62 

63 def testHtcEscapeNoSideEffect(self): 

64 val = "'val'" 

65 self.assertEqual(lssthtc.htc_escape(val), "''val''") 

66 self.assertEqual(val, "'val'") 

67 

68 def testHtcEscapeQuot(self): 

69 self.assertEqual(lssthtc.htc_escape("&quot;val&quot;"), '"val"') 

70 

71 def testHtcVersion(self): 

72 ver = lssthtc.htc_version() 

73 self.assertRegex(ver, r"^\d+\.\d+\.\d+$") 

74 

75 

76class HtcTweakJobInfoTestCase(unittest.TestCase): 

77 """Test the function responsible for massaging job information.""" 

78 

79 def setUp(self): 

80 self.log_dir = tempfile.TemporaryDirectory() 

81 self.log_dirname = pathlib.Path(self.log_dir.name) 

82 self.job = { 

83 "Cluster": 1, 

84 "Proc": 0, 

85 "Iwd": str(self.log_dirname), 

86 "Owner": self.log_dirname.owner(), 

87 "MyType": None, 

88 "TerminatedNormally": True, 

89 } 

90 

91 def tearDown(self): 

92 self.log_dir.cleanup() 

93 

94 def testDirectAssignments(self): 

95 lssthtc.htc_tweak_log_info(self.log_dirname, self.job) 

96 self.assertEqual(self.job["ClusterId"], self.job["Cluster"]) 

97 self.assertEqual(self.job["ProcId"], self.job["Proc"]) 

98 self.assertEqual(self.job["Iwd"], str(self.log_dirname)) 

99 self.assertEqual(self.job["Owner"], self.log_dirname.owner()) 

100 

101 def testIncompatibleAdPassThru(self): 

102 # Passing a job ad with insufficient information should be a no-op. 

103 expected = {"foo": "bar"} 

104 result = dict(expected) 

105 lssthtc.htc_tweak_log_info(self.log_dirname, result) 

106 self.assertEqual(result, expected) 

107 

108 def testJobStatusAssignmentJobAbortedEvent(self): 

109 job = self.job | {"MyType": "JobAbortedEvent"} 

110 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

111 self.assertTrue("JobStatus" in job) 

112 self.assertEqual(job["JobStatus"], htcondor.JobStatus.REMOVED) 

113 

114 def testJobStatusAssignmentExecuteEvent(self): 

115 job = self.job | {"MyType": "ExecuteEvent"} 

116 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

117 self.assertTrue("JobStatus" in job) 

118 self.assertEqual(job["JobStatus"], htcondor.JobStatus.RUNNING) 

119 

120 def testJobStatusAssignmentSubmitEvent(self): 

121 job = self.job | {"MyType": "SubmitEvent"} 

122 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

123 self.assertTrue("JobStatus" in job) 

124 self.assertEqual(job["JobStatus"], htcondor.JobStatus.IDLE) 

125 

126 def testJobStatusAssignmentJobHeldEvent(self): 

127 job = self.job | {"MyType": "JobHeldEvent"} 

128 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

129 self.assertTrue("JobStatus" in job) 

130 self.assertEqual(job["JobStatus"], htcondor.JobStatus.HELD) 

131 

132 def testJobStatusAssignmentJobTerminatedEvent(self): 

133 job = self.job | {"MyType": "JobTerminatedEvent"} 

134 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

135 self.assertTrue("JobStatus" in job) 

136 self.assertEqual(job["JobStatus"], htcondor.JobStatus.COMPLETED) 

137 

138 def testJobStatusAssignmentPostScriptTerminatedEvent(self): 

139 job = self.job | {"MyType": "PostScriptTerminatedEvent"} 

140 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

141 self.assertTrue("JobStatus" in job) 

142 self.assertEqual(job["JobStatus"], htcondor.JobStatus.COMPLETED) 

143 

144 def testJobStatusAssignmentReleaseEventMainDagJob(self): 

145 job = self.job | {"MyType": "JobReleaseEvent"} 

146 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

147 self.assertTrue("JobStatus" in job) 

148 self.assertEqual(job["JobStatus"], htcondor.JobStatus.RUNNING) 

149 

150 def testJobStatusAssignmentReleaseEventForNodeJob(self): 

151 job = self.job | {"MyType": "JobReleaseEvent", "DAGNodeName": "test_payload_job"} 

152 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

153 self.assertTrue("JobStatus" in job) 

154 self.assertEqual(job["JobStatus"], None) 

155 

156 def testAddingExitStatusSuccess(self): 

157 job = self.job | { 

158 "MyType": "JobTerminatedEvent", 

159 "ToE": {"ExitBySignal": False, "ExitCode": 1}, 

160 } 

161 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

162 self.assertIn("ExitBySignal", job) 

163 self.assertIs(job["ExitBySignal"], False) 

164 self.assertIn("ExitCode", job) 

165 self.assertEqual(job["ExitCode"], 1) 

166 

167 def testAddingExitStatusFailure(self): 

168 job = self.job | { 

169 "MyType": "JobHeldEvent", 

170 } 

171 with self.assertLogs(logger=logger, level="ERROR") as cm: 

172 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

173 self.assertIn("Could not determine exit status", cm.output[0]) 

174 

175 def testLoggingUnknownLogEvent(self): 

176 job = self.job | {"MyType": "Foo"} 

177 with self.assertLogs(logger=logger, level="DEBUG") as cm: 

178 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

179 self.assertIn("Unknown log event", cm.output[1]) 

180 

181 def testMissingKey(self): 

182 job = self.job 

183 del job["Cluster"] 

184 with self.assertRaises(KeyError) as cm: 

185 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

186 self.assertEqual(str(cm.exception), "'Cluster'") 

187 

188 

189class HtcCheckDagmanOutputTestCase(unittest.TestCase): 

190 """Test htc_check_dagman_output function.""" 

191 

192 def test_missing_output_file(self): 

193 with temporaryDirectory() as tmp_dir: 

194 with self.assertRaises(FileNotFoundError): 

195 _ = lssthtc.htc_check_dagman_output(tmp_dir) 

196 

197 def test_permissions_output_file(self): 

198 with temporaryDirectory() as tmp_dir: 

199 copy2(f"{TESTDIR}/data/test_tmpdir_abort.dag.dagman.out", tmp_dir) 

200 os.chmod(f"{tmp_dir}/test_tmpdir_abort.dag.dagman.out", 0o200) 

201 print(os.stat(f"{tmp_dir}/test_tmpdir_abort.dag.dagman.out")) 

202 results = lssthtc.htc_check_dagman_output(tmp_dir) 

203 os.chmod(f"{tmp_dir}/test_tmpdir_abort.dag.dagman.out", 0o600) 

204 self.assertIn("Could not read dagman output file", results) 

205 

206 def test_submit_failure(self): 

207 with temporaryDirectory() as tmp_dir: 

208 copy2(f"{TESTDIR}/data/bad_submit.dag.dagman.out", tmp_dir) 

209 results = lssthtc.htc_check_dagman_output(tmp_dir) 

210 self.assertIn("Warn: Job submission issues (last: ", results) 

211 

212 def test_tmpdir_abort(self): 

213 with temporaryDirectory() as tmp_dir: 

214 copy2(f"{TESTDIR}/data/test_tmpdir_abort.dag.dagman.out", tmp_dir) 

215 results = lssthtc.htc_check_dagman_output(tmp_dir) 

216 self.assertIn("Cannot submit from /tmp", results) 

217 

218 def test_no_messages(self): 

219 with temporaryDirectory() as tmp_dir: 

220 copy2(f"{TESTDIR}/data/test_no_messages.dag.dagman.out", tmp_dir) 

221 results = lssthtc.htc_check_dagman_output(tmp_dir) 

222 self.assertEqual("", results) 

223 

224 

225class SummarizeDagTestCase(unittest.TestCase): 

226 """Test summarize_dag function.""" 

227 

228 def test_no_dag_file(self): 

229 with temporaryDirectory() as tmp_dir: 

230 summary, job_name_to_pipetask, job_name_to_type = lssthtc.summarize_dag(tmp_dir) 

231 self.assertFalse(len(job_name_to_pipetask)) 

232 self.assertFalse(len(job_name_to_type)) 

233 self.assertFalse(summary) 

234 

235 def test_success(self): 

236 with temporaryDirectory() as tmp_dir: 

237 copy2(f"{TESTDIR}/data/good.dag", tmp_dir) 

238 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir) 

239 self.assertEqual(summary, "pipetaskInit:1;label1:1;label2:1;label3:1;finalJob:1") 

240 self.assertEqual( 

241 job_name_to_label, 

242 { 

243 "pipetaskInit": "pipetaskInit", 

244 "0682f8f9-12f0-40a5-971e-8b30c7231e5c_label1_val1_val2": "label1", 

245 "d0305e2d-f164-4a85-bd24-06afe6c84ed9_label2_val1_val2": "label2", 

246 "2806ecc9-1bba-4362-8fff-ab4e6abb9f83_label3_val1_val2": "label3", 

247 "finalJob": "finalJob", 

248 }, 

249 ) 

250 self.assertEqual( 

251 job_name_to_type, 

252 { 

253 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

254 "0682f8f9-12f0-40a5-971e-8b30c7231e5c_label1_val1_val2": lssthtc.WmsNodeType.PAYLOAD, 

255 "d0305e2d-f164-4a85-bd24-06afe6c84ed9_label2_val1_val2": lssthtc.WmsNodeType.PAYLOAD, 

256 "2806ecc9-1bba-4362-8fff-ab4e6abb9f83_label3_val1_val2": lssthtc.WmsNodeType.PAYLOAD, 

257 "finalJob": lssthtc.WmsNodeType.FINAL, 

258 }, 

259 ) 

260 

261 def test_service(self): 

262 with temporaryDirectory() as tmp_dir: 

263 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag", tmp_dir) 

264 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir) 

265 self.assertEqual(summary, "pipetaskInit:1;label1:2;label2:2;finalJob:1") 

266 self.assertEqual( 

267 job_name_to_label, 

268 { 

269 "pipetaskInit": "pipetaskInit", 

270 "057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b": "label1", 

271 "4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a": "label1", 

272 "40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b": "label2", 

273 "696ee50d-e711-40d6-9caf-ee29ae4a656d_label2_val1a_val2a": "label2", 

274 "finalJob": "finalJob", 

275 "provisioningJob": "provisioningJob", 

276 }, 

277 ) 

278 self.assertEqual( 

279 job_name_to_type, 

280 { 

281 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

282 "057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

283 "4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

284 "40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

285 "696ee50d-e711-40d6-9caf-ee29ae4a656d_label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

286 "finalJob": lssthtc.WmsNodeType.FINAL, 

287 "provisioningJob": lssthtc.WmsNodeType.SERVICE, 

288 }, 

289 ) 

290 

291 def test_noop(self): 

292 with temporaryDirectory() as tmp_dir: 

293 copy2(f"{TESTDIR}/data/noop_running_1/noop_running_1.dag", tmp_dir) 

294 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir) 

295 self.assertEqual( 

296 set(summary.split(";")), 

297 {"pipetaskInit:1", "label1:6", "label2:6", "label3:6", "label4:6", "label5:6", "finalJob:1"}, 

298 ) 

299 self.assertEqual( 

300 job_name_to_label, 

301 { 

302 "label1_val1a_val2a": "label1", 

303 "label1_val1a_val2b": "label1", 

304 "label1_val1b_val2a": "label1", 

305 "label1_val1b_val2b": "label1", 

306 "label1_val1c_val2a": "label1", 

307 "label1_val1c_val2b": "label1", 

308 "label2_val1a_val2a": "label2", 

309 "label2_val1a_val2b": "label2", 

310 "label2_val1b_val2a": "label2", 

311 "label2_val1b_val2b": "label2", 

312 "label2_val1c_val2a": "label2", 

313 "label2_val1c_val2b": "label2", 

314 "label3_val1a_val2a": "label3", 

315 "label3_val1a_val2b": "label3", 

316 "label3_val1b_val2a": "label3", 

317 "label3_val1b_val2b": "label3", 

318 "label3_val1c_val2a": "label3", 

319 "label3_val1c_val2b": "label3", 

320 "label4_val1a_val2a": "label4", 

321 "label4_val1a_val2b": "label4", 

322 "label4_val1b_val2a": "label4", 

323 "label4_val1b_val2b": "label4", 

324 "label4_val1c_val2a": "label4", 

325 "label4_val1c_val2b": "label4", 

326 "label5_val1a_val2a": "label5", 

327 "label5_val1a_val2b": "label5", 

328 "label5_val1b_val2a": "label5", 

329 "label5_val1b_val2b": "label5", 

330 "label5_val1c_val2a": "label5", 

331 "label5_val1c_val2b": "label5", 

332 "finalJob": "finalJob", 

333 "pipetaskInit": "pipetaskInit", 

334 "wms_noop_order1_val1a": "order1", 

335 "wms_noop_order1_val1b": "order1", 

336 }, 

337 ) 

338 self.assertEqual( 

339 job_name_to_type, 

340 { 

341 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

342 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

343 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

344 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

345 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

346 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

347 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

348 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

349 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

350 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

351 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

352 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

353 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

354 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

355 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

356 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

357 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

358 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

359 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

360 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

361 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

362 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

363 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

364 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

365 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

366 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

367 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

368 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

369 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

370 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

371 "finalJob": lssthtc.WmsNodeType.FINAL, 

372 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

373 "wms_noop_order1_val1a": lssthtc.WmsNodeType.NOOP, 

374 "wms_noop_order1_val1b": lssthtc.WmsNodeType.NOOP, 

375 }, 

376 ) 

377 

378 def test_subdags(self): 

379 self.maxDiff = None 

380 with temporaryDirectory() as tmp_dir: 

381 submit_dir = os.path.join(tmp_dir, "group_running_1") 

382 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

383 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(submit_dir) 

384 self.assertEqual( 

385 set(summary.split(";")), 

386 {"pipetaskInit:1", "label1:6", "label2:6", "label3:6", "label4:6", "label5:6", "finalJob:1"}, 

387 ) 

388 

389 self.assertEqual( 

390 job_name_to_label, 

391 { 

392 "pipetaskInit": "pipetaskInit", 

393 "label1_val1b_val2a": "label1", 

394 "label1_val1c_val2a": "label1", 

395 "label1_val1a_val2b": "label1", 

396 "label1_val1b_val2b": "label1", 

397 "label1_val1c_val2b": "label1", 

398 "label1_val1a_val2a": "label1", 

399 "label2_val1a_val2b": "label2", 

400 "label2_val1a_val2a": "label2", 

401 "label2_val1b_val2a": "label2", 

402 "label2_val1b_val2b": "label2", 

403 "label2_val1c_val2a": "label2", 

404 "label2_val1c_val2b": "label2", 

405 "label3_val1b_val2a": "label3", 

406 "label3_val1c_val2a": "label3", 

407 "label3_val1a_val2b": "label3", 

408 "label3_val1b_val2b": "label3", 

409 "label3_val1c_val2b": "label3", 

410 "label3_val1a_val2a": "label3", 

411 "label4_val1a_val2b": "label4", 

412 "label4_val1a_val2a": "label4", 

413 "label4_val1b_val2a": "label4", 

414 "label4_val1b_val2b": "label4", 

415 "label4_val1c_val2a": "label4", 

416 "label4_val1c_val2b": "label4", 

417 "label5_val1a_val2b": "label5", 

418 "label5_val1a_val2a": "label5", 

419 "label5_val1b_val2a": "label5", 

420 "label5_val1b_val2b": "label5", 

421 "label5_val1c_val2a": "label5", 

422 "label5_val1c_val2b": "label5", 

423 "finalJob": "finalJob", 

424 "provisioningJob": "provisioningJob", 

425 "wms_group_order1_val1a": "order1", 

426 "wms_group_order1_val1b": "order1", 

427 "wms_group_order1_val1c": "order1", 

428 "wms_check_status_wms_group_order1_val1a": "order1", 

429 "wms_check_status_wms_group_order1_val1b": "order1", 

430 "wms_check_status_wms_group_order1_val1c": "order1", 

431 }, 

432 ) 

433 

434 self.assertEqual( 

435 job_name_to_type, 

436 { 

437 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

438 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

439 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

440 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

441 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

442 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

443 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

444 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

445 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

446 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

447 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

448 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

449 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

450 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

451 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

452 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

453 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

454 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

455 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

456 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

457 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

458 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

459 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

460 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

461 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

462 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

463 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

464 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

465 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

466 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

467 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

468 "finalJob": lssthtc.WmsNodeType.FINAL, 

469 "provisioningJob": lssthtc.WmsNodeType.SERVICE, 

470 "wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG, 

471 "wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG, 

472 "wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG, 

473 "wms_check_status_wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG_CHECK, 

474 "wms_check_status_wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG_CHECK, 

475 "wms_check_status_wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG_CHECK, 

476 }, 

477 ) 

478 

479 

480class ReadDagNodesLogTestCase(unittest.TestCase): 

481 """Test read_dag_nodes_log function.""" 

482 

483 def setUp(self): 

484 self.tmpdir = tempfile.mkdtemp() 

485 

486 def tearDown(self): 

487 rmtree(self.tmpdir, ignore_errors=True) 

488 

489 def testFileMissing(self): 

490 with self.assertRaisesRegex(FileNotFoundError, "DAGMan node log not found in"): 

491 _ = lssthtc.read_dag_nodes_log(self.tmpdir) 

492 

493 def testRegular(self): 

494 with temporaryDirectory() as tmp_dir: 

495 submit_dir = os.path.join(tmp_dir, "tiny_problems") 

496 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

497 results = lssthtc.read_dag_nodes_log(submit_dir) 

498 self.assertEqual(results["9231.0"]["Cluster"], 9231) 

499 self.assertEqual(results["9231.0"]["Proc"], 0) 

500 self.assertEqual(results["9231.0"]["ToE"]["ExitCode"], 1) 

501 self.assertEqual(len(results), 6) 

502 

503 def testSubdags(self): 

504 """Making sure it gets data from subdag dirs and doesn't 

505 fail if some subdags haven't started running yet. 

506 """ 

507 with temporaryDirectory() as tmp_dir: 

508 submit_dir = os.path.join(tmp_dir, "group_running_1") 

509 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

510 results = lssthtc.read_dag_nodes_log(submit_dir) 

511 # main dag 

512 self.assertEqual(results["10094.0"]["Cluster"], 10094) 

513 # subdag 

514 self.assertEqual(results["10112.0"]["Cluster"], 10112) 

515 self.assertEqual(results["10116.0"]["Cluster"], 10116) 

516 

517 

518class ReadNodeStatusTestCase(unittest.TestCase): 

519 """Test read_node_status function.""" 

520 

521 def setUp(self): 

522 self.tmpdir = tempfile.mkdtemp() 

523 

524 def tearDown(self): 

525 rmtree(self.tmpdir, ignore_errors=True) 

526 

527 def testServiceJobNotSubmitted(self): 

528 # tiny_prov_no_submit files have successful workflow 

529 # but provisioningJob could not submit. 

530 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag.nodes.log", self.tmpdir) 

531 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag.dagman.log", self.tmpdir) 

532 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.node_status", self.tmpdir) 

533 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag", self.tmpdir) 

534 

535 jobs = lssthtc.read_node_status(self.tmpdir) 

536 found = [ 

537 id_ 

538 for id_ in jobs 

539 if jobs[id_].get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN) == lssthtc.WmsNodeType.SERVICE 

540 ] 

541 self.assertEqual(len(found), 1) 

542 self.assertEqual(jobs[found[0]]["DAGNodeName"], "provisioningJob") 

543 self.assertEqual(jobs[found[0]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY) 

544 

545 def testMissingStatusFile(self): 

546 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag.nodes.log", self.tmpdir) 

547 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag.dagman.log", self.tmpdir) 

548 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag", self.tmpdir) 

549 

550 jobs = lssthtc.read_node_status(self.tmpdir) 

551 self.assertEqual(len(jobs), 7) 

552 self.assertEqual(jobs["9230.0"]["DAGNodeName"], "pipetaskInit") 

553 self.assertEqual(jobs["9230.0"]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD) 

554 found = [ 

555 id_ 

556 for id_ in jobs 

557 if jobs[id_].get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN) == lssthtc.WmsNodeType.SERVICE 

558 ] 

559 self.assertEqual(len(found), 1) 

560 self.assertEqual(jobs[found[0]]["DAGNodeName"], "provisioningJob") 

561 

562 def testSubdagsRunning(self): 

563 with temporaryDirectory() as tmp_dir: 

564 test_tmp_dir = pathlib.Path(tmp_dir) 

565 submit_dir = test_tmp_dir / "submit" 

566 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

567 jobs = lssthtc.read_node_status(submit_dir) 

568 self.assertEqual(len(jobs), 39) # includes non-payload jobs 

569 # not guaranteed ids are same, so use names instead 

570 job_name_to_id = {} 

571 for id_, info in jobs.items(): 

572 job_name_to_id[info.get("DAGNodeName", id_)] = id_ 

573 job_type_to_names = {} 

574 for id_, info in jobs.items(): 

575 job_type_to_names.setdefault( 

576 info.get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN), set() 

577 ).add(info.get("DAGNodeName", id_)) 

578 

579 # check counts 

580 self.assertNotIn(lssthtc.WmsNodeType.NOOP, job_type_to_names) 

581 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.PAYLOAD]), 31) 

582 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.FINAL]), 1) 

583 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SERVICE]), 1) 

584 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG]), 3) 

585 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG_CHECK]), 3) 

586 

587 # spot check some statuses 

588 self.assertEqual( 

589 jobs[job_name_to_id["label3_val1a_val2b"]]["NodeStatus"], lssthtc.NodeStatus.DONE 

590 ) 

591 self.assertEqual( 

592 jobs[job_name_to_id["wms_group_order1_val1a"]]["NodeStatus"], lssthtc.NodeStatus.SUBMITTED 

593 ) 

594 self.assertEqual( 

595 jobs[job_name_to_id["label5_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY 

596 ) 

597 self.assertEqual( 

598 jobs[job_name_to_id["label2_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.DONE 

599 ) 

600 

601 def testSubdagsFailed(self): 

602 with temporaryDirectory() as tmp_dir: 

603 test_tmp_dir = pathlib.Path(tmp_dir) 

604 submit_dir = test_tmp_dir / "submit" 

605 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

606 jobs = lssthtc.read_node_status(submit_dir) 

607 self.assertEqual(len(jobs), 39) 

608 # not guaranteed ids are same, so use names instead 

609 job_name_to_id = {} 

610 for id_, info in jobs.items(): 

611 job_name_to_id[info.get("DAGNodeName", id_)] = id_ 

612 job_type_to_names = {} 

613 for id_, info in jobs.items(): 

614 job_type_to_names.setdefault( 

615 info.get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN), set() 

616 ).add(info.get("DAGNodeName", id_)) 

617 

618 # check counts 

619 self.assertNotIn(lssthtc.WmsNodeType.NOOP, job_type_to_names) 

620 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.PAYLOAD]), 31) 

621 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.FINAL]), 1) 

622 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SERVICE]), 1) 

623 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG]), 3) 

624 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG_CHECK]), 3) 

625 

626 # spot check some statuses 

627 self.assertEqual( 

628 jobs[job_name_to_id["label3_val1a_val2b"]]["NodeStatus"], lssthtc.NodeStatus.DONE 

629 ) 

630 self.assertEqual( 

631 jobs[job_name_to_id["wms_group_order1_val1a"]]["NodeStatus"], lssthtc.NodeStatus.DONE 

632 ) 

633 self.assertEqual( 

634 jobs[job_name_to_id["label5_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.DONE 

635 ) 

636 

637 self.assertEqual( 

638 jobs[job_name_to_id["label5_val1b_val2a"]]["NodeStatus"], lssthtc.NodeStatus.FUTILE 

639 ) 

640 self.assertEqual( 

641 jobs[job_name_to_id["wms_group_order1_val1b"]]["NodeStatus"], lssthtc.NodeStatus.DONE 

642 ) 

643 self.assertEqual( 

644 jobs[job_name_to_id["wms_check_status_wms_group_order1_val1b"]]["NodeStatus"], 

645 lssthtc.NodeStatus.ERROR, 

646 ) 

647 

648 

649class ReadSingleNodeStatusTestCase(unittest.TestCase): 

650 """Test read_single_node_status function.""" 

651 

652 def setUp(self): 

653 self.tmpdir = tempfile.mkdtemp() 

654 

655 def tearDown(self): 

656 rmtree(self.tmpdir, ignore_errors=True) 

657 

658 def _copyFiles(self, data_subdir, suffixes): 

659 """Copy files with given suffixes from tests/data/<data_subdir>/.""" 

660 for suffix in suffixes: 

661 copy2(f"{TESTDIR}/data/{data_subdir}/{data_subdir}{suffix}", self.tmpdir) 

662 

663 def _jobNameToId(self, jobs): 

664 return {info["DAGNodeName"]: id_ for id_, info in jobs.items()} 

665 

666 def testAllDone(self): 

667 self._copyFiles( 

668 "tiny_success", 

669 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"], 

670 ) 

671 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status" 

672 jobs = lssthtc.read_single_node_status(filename, -1) 

673 

674 self.assertEqual(len(jobs), 5) 

675 name_to_id = self._jobNameToId(jobs) 

676 

677 # All four submitted nodes are marked DONE. 

678 for name in [ 

679 "pipetaskInit", 

680 "5bba27bd-8df7-4668-a9c5-e911192c5cdb_label1_val1_val2", 

681 "0b225f1f-6edf-4380-b546-76c97947a88f_label2_val1_val2", 

682 "finalJob", 

683 ]: 

684 self.assertIn(name, name_to_id, msg=f"Missing job {name}") 

685 self.assertEqual( 

686 jobs[name_to_id[name]]["NodeStatus"], 

687 lssthtc.NodeStatus.DONE, 

688 msg=f"Expected DONE for {name}", 

689 ) 

690 

691 # Service job not tracked by node_status; it came from the event log so 

692 # it has a real positive ClusterId but no NodeStatus field. 

693 self.assertIn("provisioningJob", name_to_id) 

694 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0) 

695 

696 # Spot-check labels and types. 

697 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["bps_job_label"], "pipetaskInit") 

698 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD) 

699 self.assertEqual(jobs[name_to_id["finalJob"]]["wms_node_type"], lssthtc.WmsNodeType.FINAL) 

700 self.assertEqual(jobs[name_to_id["provisioningJob"]]["wms_node_type"], lssthtc.WmsNodeType.SERVICE) 

701 

702 # DAGManJobID is populated from the dagman log for every job. 

703 for job in jobs.values(): 

704 self.assertIn("DAGManJobID", job) 

705 

706 def testMixedStatuses(self): 

707 self._copyFiles( 

708 "tiny_problems", 

709 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"], 

710 ) 

711 filename = pathlib.Path(self.tmpdir) / "tiny_problems.node_status" 

712 jobs = lssthtc.read_single_node_status(filename, -1) 

713 

714 self.assertEqual(len(jobs), 7) 

715 name_to_id = self._jobNameToId(jobs) 

716 

717 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE) 

718 self.assertEqual( 

719 jobs[name_to_id["057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b"]]["NodeStatus"], 

720 lssthtc.NodeStatus.ERROR, 

721 ) 

722 self.assertEqual( 

723 jobs[name_to_id["4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a"]]["NodeStatus"], 

724 lssthtc.NodeStatus.DONE, 

725 ) 

726 self.assertEqual( 

727 jobs[name_to_id["40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b"]]["NodeStatus"], 

728 lssthtc.NodeStatus.FUTILE, 

729 ) 

730 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.ERROR) 

731 # Service job not tracked by node_status; came from event log so 

732 # it has a real positive ClusterId but no NodeStatus field. 

733 self.assertIn("provisioningJob", name_to_id) 

734 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0) 

735 

736 def testRunningWorkflow(self): 

737 self._copyFiles( 

738 "tiny_running", 

739 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"], 

740 ) 

741 filename = pathlib.Path(self.tmpdir) / "tiny_running.node_status" 

742 jobs = lssthtc.read_single_node_status(filename, -1) 

743 

744 self.assertEqual(len(jobs), 5) 

745 name_to_id = self._jobNameToId(jobs) 

746 

747 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE) 

748 self.assertEqual( 

749 jobs[name_to_id["ca27ea57-c014-44c1-838a-78c06bc3ec1b_label1_val1_val2"]]["NodeStatus"], 

750 lssthtc.NodeStatus.SUBMITTED, 

751 ) 

752 self.assertEqual( 

753 jobs[name_to_id["dbf919fa-5453-4b05-8806-ad6390fda0a3_label2_val1_val2"]]["NodeStatus"], 

754 lssthtc.NodeStatus.NOT_READY, 

755 ) 

756 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY) 

757 # Service job appeared in the event log; has a real positive ClusterId. 

758 self.assertIn("provisioningJob", name_to_id) 

759 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0) 

760 

761 def testMissingNodeStatusFile(self): 

762 # Omit the .node_status file; jobs must be built from the event log 

763 # and dag. 

764 self._copyFiles( 

765 "tiny_problems", 

766 [".dag", ".dag.dagman.log", ".dag.nodes.log"], 

767 ) 

768 filename = pathlib.Path(self.tmpdir) / "tiny_problems.node_status" 

769 jobs = lssthtc.read_single_node_status(filename, -1) 

770 

771 self.assertEqual(len(jobs), 7) 

772 name_to_id = self._jobNameToId(jobs) 

773 

774 # Jobs that appeared in the event log have real (positive) cluster IDs. 

775 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["DAGNodeName"], "pipetaskInit") 

776 self.assertGreater(jobs[name_to_id["pipetaskInit"]]["ClusterId"], 0) 

777 

778 # The service job appeared in the event log and has a real positive ID. 

779 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0) 

780 

781 # All jobs carry the correct label and type from the dag file. 

782 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD) 

783 self.assertEqual(jobs[name_to_id["provisioningJob"]]["wms_node_type"], lssthtc.WmsNodeType.SERVICE) 

784 

785 def testMissingLogFiles(self): 

786 # Omit both log files; every job should get a fake negative ID. 

787 self._copyFiles("tiny_success", [".dag", ".node_status"]) 

788 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status" 

789 jobs = lssthtc.read_single_node_status(filename, -1) 

790 

791 self.assertEqual(len(jobs), 5) 

792 for job in jobs.values(): 

793 self.assertLess(job["ClusterId"], 0) 

794 

795 # NodeStatus values from the node_status file must still be preserved. 

796 name_to_id = self._jobNameToId(jobs) 

797 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE) 

798 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.DONE) 

799 self.assertEqual(jobs[name_to_id["provisioningJob"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY) 

800 

801 def testInitFakeId(self): 

802 # Verify fake IDs count down from the given starting value. 

803 self._copyFiles("tiny_success", [".dag", ".node_status"]) 

804 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status" 

805 init_fake_id = -10 

806 jobs = lssthtc.read_single_node_status(filename, init_fake_id) 

807 

808 self.assertEqual(len(jobs), 5) 

809 cluster_ids = [job["ClusterId"] for job in jobs.values()] 

810 # All IDs must be at most init_fake_id (i.e., -10 or lower). 

811 for cid in cluster_ids: 

812 self.assertLessEqual(cid, init_fake_id) 

813 # All IDs must be unique. 

814 self.assertEqual(len(set(cluster_ids)), len(cluster_ids)) 

815 

816 def testFromDagJobAttribute(self): 

817 self._copyFiles( 

818 "tiny_success", 

819 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"], 

820 ) 

821 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status" 

822 jobs = lssthtc.read_single_node_status(filename, -1) 

823 

824 for id_, job in jobs.items(): 

825 self.assertEqual( 

826 job["from_dag_job"], 

827 "wms_tiny_success", 

828 msg=f"Job {id_} has wrong from_dag_job", 

829 ) 

830 

831 def testServiceJobPlaceholder(self): 

832 self._copyFiles( 

833 "tiny_prov_no_submit", 

834 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"], 

835 ) 

836 filename = pathlib.Path(self.tmpdir) / "tiny_prov_no_submit.node_status" 

837 jobs = lssthtc.read_single_node_status(filename, -1) 

838 

839 service_jobs = [ 

840 (id_, info) 

841 for id_, info in jobs.items() 

842 if info.get("wms_node_type") == lssthtc.WmsNodeType.SERVICE 

843 ] 

844 self.assertEqual(len(service_jobs), 1) 

845 service_id, service_job = service_jobs[0] 

846 self.assertEqual(service_job["DAGNodeName"], "provisioningJob") 

847 self.assertEqual(service_job["NodeStatus"], lssthtc.NodeStatus.NOT_READY) 

848 self.assertLess(service_job["ClusterId"], 0) 

849 

850 

851class HTCJobTestCase(unittest.TestCase): 

852 """Test HTCJob methods.""" 

853 

854 def testWriteDagCommandsPayload(self): 

855 job = lssthtc.HTCJob( 

856 "job1", 

857 "label1", 

858 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"}, 

859 {"dir": "jobs/label1"}, 

860 ) 

861 job.subfile = "job1.sub" 

862 

863 mockfh = io.StringIO() 

864 job.write_dag_commands(mockfh, "../..") 

865 self.assertIn('JOB job1 "job1.sub" DIR "../../jobs/label1"', mockfh.getvalue()) 

866 

867 def testWriteDagCommandsNotJob(self): 

868 # Testing giving command_name, no dag_rel_path and no dir 

869 job = lssthtc.HTCJob( 

870 "finalJob", 

871 "finalJob", 

872 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"}, 

873 ) 

874 job.subfile = "jobs/finalJob/finalJob.sub" 

875 mockfh = io.StringIO() 

876 job.write_dag_commands(mockfh, "", "FINAL") 

877 self.assertIn('FINAL finalJob "jobs/finalJob/finalJob.sub"', mockfh.getvalue()) 

878 

879 def testWriteDagCommandsNoop(self): 

880 job = lssthtc.HTCJob("wms_noop_job1", "label1", {}, {"noop": True}) 

881 job.subfile = "notthere.sub" 

882 mockfh = io.StringIO() 

883 job.write_dag_commands(mockfh, "") 

884 self.assertIn("NOOP", mockfh.getvalue()) 

885 

886 def testWriteSubmitFile(self): 

887 job = lssthtc.HTCJob( 

888 "job1", 

889 "label1", 

890 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"}, 

891 ) 

892 with temporaryDirectory() as tmp_dir: 

893 filename = pathlib.Path(tmp_dir) / "label1/job1.sub" 

894 job.write_submit_file(filename.parent) 

895 self.assertTrue(filename.exists()) 

896 # Try to make Submit object from file to find any syntax issues 

897 _ = lssthtc.htc_create_submit_from_file(filename) 

898 

899 def testWriteSubmitFileExists(self): 

900 job = lssthtc.HTCJob( 

901 "job1", 

902 "label1", 

903 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"}, 

904 ) 

905 with temporaryDirectory() as tmp_dir: 

906 filename = pathlib.Path(tmp_dir) / "job1.sub" 

907 job.subfile = filename 

908 with open(filename, "w"): 

909 pass # make empty file 

910 job.write_submit_file(filename.parent) 

911 # make sure didn't overwrite file 

912 self.assertEqual(filename.stat().st_size, 0, "Incorrectly overwrote existing file") 

913 

914 

915class HtcWriteJobCommands(unittest.TestCase): 

916 """Test _htc_write_job_commands function.""" 

917 

918 def testAllCommands(self): 

919 dag_cmds = { 

920 "pre": { 

921 "defer": {"status": 1, "time": 120}, 

922 "debug": {"filename": "debug_pre.txt", "type": "ALL"}, 

923 "executable": "exec1", 

924 "arguments": "arg1 arg2", 

925 }, 

926 "post": { 

927 "defer": {"status": 2, "time": 180}, 

928 "debug": {"filename": "debug_post.txt", "type": "ALL"}, 

929 "executable": "exec2", 

930 "arguments": "arg3 arg4", 

931 }, 

932 "vars": {"num": 8, "spaces": "a space"}, 

933 "pre_skip": "1", 

934 "retry": 3, 

935 "retry_unless_exit": 1, 

936 "abort_dag_on": {"node_exit": 100, "abort_exit": 4}, 

937 "priority": 123, 

938 } 

939 

940 truth = """SCRIPT DEFER 1 120 DEBUG debug_pre.txt ALL PRE job1 exec1 arg1 arg2 

941SCRIPT DEFER 2 180 DEBUG debug_post.txt ALL POST job1 exec2 arg3 arg4 

942VARS job1 num="8" 

943VARS job1 spaces="a space" 

944PRE_SKIP job1 1 

945RETRY job1 3 UNLESS-EXIT 1 

946ABORT-DAG-ON job1 100 RETURN 4 

947PRIORITY job1 123 

948""" 

949 mockfh = io.StringIO() 

950 lssthtc._htc_write_job_commands(mockfh, "job1", dag_cmds) 

951 self.assertEqual(mockfh.getvalue(), truth) 

952 

953 def testPartialCommands(self): 

954 # Trigger skipping the inner if clauses. 

955 dag_cmds = { 

956 "pre": { 

957 "executable": "exec1", 

958 }, 

959 "post": { 

960 "executable": "exec2", 

961 }, 

962 "vars": {"num": 8, "spaces": "a space"}, 

963 "pre_skip": "1", 

964 "retry": 3, 

965 } 

966 

967 truth = """SCRIPT PRE job1 exec1 

968SCRIPT POST job1 exec2 

969VARS job1 num="8" 

970VARS job1 spaces="a space" 

971PRE_SKIP job1 1 

972RETRY job1 3 

973""" 

974 mockfh = io.StringIO() 

975 lssthtc._htc_write_job_commands(mockfh, "job1", dag_cmds) 

976 self.assertEqual(mockfh.getvalue(), truth) 

977 

978 def testNoCommands(self): 

979 dag_cmds = {} 

980 mockfh = io.StringIO() 

981 lssthtc._htc_write_job_commands(mockfh, "job2", dag_cmds) 

982 self.assertEqual(mockfh.getvalue(), "") 

983 

984 def testFinal(self): 

985 self.maxDiff = None 

986 dag_cmds = { 

987 "pre": { 

988 "defer": {"status": 1, "time": 120}, 

989 "debug": {"filename": "debug_pre.txt", "type": "ALL"}, 

990 "executable": "exec1", 

991 "arguments": "arg1 arg2", 

992 }, 

993 "post": { 

994 "defer": {"status": 2, "time": 180}, 

995 "debug": {"filename": "debug_post.txt", "type": "ALL"}, 

996 "executable": "exec2", 

997 "arguments": "arg3 arg4", 

998 }, 

999 "vars": {"num": 8, "spaces": "a space"}, 

1000 "pre_skip": "1", 

1001 "retry": 3, 

1002 "retry_unless_exit": 1, 

1003 "abort_dag_on": {"node_exit": 100, "abort_exit": 4}, 

1004 "priority": 123, 

1005 } 

1006 

1007 truth = """SCRIPT DEFER 1 120 DEBUG debug_pre.txt ALL PRE finalJob exec1 arg1 arg2 

1008SCRIPT DEFER 2 180 DEBUG debug_post.txt ALL POST finalJob exec2 arg3 arg4 

1009VARS finalJob num="8" 

1010VARS finalJob spaces="a space" 

1011PRE_SKIP finalJob 1 

1012""" 

1013 mockfh = io.StringIO() 

1014 lssthtc._htc_write_job_commands(mockfh, "finalJob", dag_cmds, "FINAL") 

1015 self.assertEqual(mockfh.getvalue(), truth) 

1016 

1017 

1018class HTCBackupFilesSinglePathTestCase(unittest.TestCase): 

1019 """Test htc_backup_files_single_path function.""" 

1020 

1021 def testSrcDestSame(self): 

1022 with temporaryDirectory() as tmp_dir: 

1023 with self.assertRaisesRegex( 

1024 RuntimeError, "Destination directory is same as the source directory" 

1025 ): 

1026 lssthtc.htc_backup_files_single_path(tmp_dir, tmp_dir) 

1027 

1028 def testSuccess(self): 

1029 with temporaryDirectory() as tmp_dir: 

1030 test_tmp_dir = pathlib.Path(tmp_dir) 

1031 submit_dir = test_tmp_dir / "the_src_dir" 

1032 copytree(f"{TESTDIR}/data/tiny_success", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1033 backup_dir = test_tmp_dir / "the_dest_dir" 

1034 backup_dir.mkdir() 

1035 lssthtc.htc_backup_files_single_path(submit_dir, backup_dir) 

1036 result_submit = [] 

1037 for root, _, files in os.walk(submit_dir): 

1038 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files]) 

1039 self.assertEqual( 

1040 set(result_submit), 

1041 { 

1042 "./tiny_success.dag.dagman.log", 

1043 "./tiny_success.dag.dagman.out", 

1044 "./tiny_success.dag", 

1045 }, 

1046 ) 

1047 result_backup = [] 

1048 for root, _, files in os.walk(backup_dir): 

1049 result_backup.extend([str(os.path.join(os.path.relpath(root, backup_dir), f)) for f in files]) 

1050 self.assertEqual( 

1051 set(result_backup), 

1052 { 

1053 "./tiny_success.info.json", 

1054 "./tiny_success.dag.metrics", 

1055 "./tiny_success.dag.nodes.log", 

1056 "./tiny_success.node_status", 

1057 }, 

1058 ) 

1059 

1060 

1061class HTCBackupFilesTestCase(unittest.TestCase): 

1062 """Test htc_backup_files function.""" 

1063 

1064 def testDirectoryNotFound(self): 

1065 with temporaryDirectory() as tmp_dir: 

1066 test_tmp_dir = pathlib.Path(tmp_dir) 

1067 submit_dir = test_tmp_dir / "submit" 

1068 with self.assertRaises(FileNotFoundError): 

1069 lssthtc.htc_backup_files(submit_dir) 

1070 

1071 def testSuccess(self): 

1072 with temporaryDirectory() as tmp_dir: 

1073 test_tmp_dir = pathlib.Path(tmp_dir) 

1074 submit_dir = test_tmp_dir / "submit" 

1075 copytree(f"{TESTDIR}/data/tiny_success", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1076 lssthtc.htc_backup_files(submit_dir) 

1077 result_submit = [] 

1078 for root, _, files in os.walk(submit_dir): 

1079 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files]) 

1080 self.assertEqual( 

1081 set(result_submit), 

1082 { 

1083 "./tiny_success.dag.dagman.log", 

1084 "./tiny_success.dag.dagman.out", 

1085 "./tiny_success.dag", 

1086 "000/tiny_success.info.json", 

1087 "000/tiny_success.dag.metrics", 

1088 "000/tiny_success.dag.nodes.log", 

1089 "000/tiny_success.node_status", 

1090 }, 

1091 ) 

1092 

1093 def testDestNotInSubmitDir(self): 

1094 with temporaryDirectory() as tmp_dir: 

1095 test_tmp_dir = pathlib.Path(tmp_dir) 

1096 submit_dir = test_tmp_dir / "submit" 

1097 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1098 with self.assertLogs("lsst.ctrl.bps.htcondor", level="WARNING") as cm: 

1099 lssthtc.htc_backup_files(submit_dir, test_tmp_dir / "backup") 

1100 self.assertIn("Invalid backup location:", cm.output[-1]) 

1101 lssthtc.htc_backup_files(submit_dir) 

1102 result_submit = [] 

1103 for root, _, files in os.walk(submit_dir): 

1104 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files]) 

1105 self.assertEqual( 

1106 set(result_submit), 

1107 { 

1108 "./tiny_problems.dag.dagman.log", 

1109 "./tiny_problems.dag.dagman.out", 

1110 "./tiny_problems.dag", 

1111 "./tiny_problems.dag.rescue001", 

1112 "001/tiny_problems.info.json", 

1113 "001/tiny_problems.dag.metrics", 

1114 "001/tiny_problems.dag.nodes.log", 

1115 "001/tiny_problems.node_status", 

1116 }, 

1117 ) 

1118 

1119 def testDestInSubmitDir(self): 

1120 with temporaryDirectory() as tmp_dir: 

1121 test_tmp_dir = pathlib.Path(tmp_dir) 

1122 submit_dir = test_tmp_dir / "submit" 

1123 backup_dir = submit_dir / "subdir" 

1124 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1125 lssthtc.htc_backup_files(submit_dir, backup_dir) 

1126 result_submit = [] 

1127 for root, _, files in os.walk(submit_dir): 

1128 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files]) 

1129 self.assertEqual( 

1130 set(result_submit), 

1131 { 

1132 "./tiny_problems.dag.dagman.log", 

1133 "./tiny_problems.dag.dagman.out", 

1134 "./tiny_problems.dag", 

1135 "./tiny_problems.dag.rescue001", 

1136 "subdir/001/tiny_problems.info.json", 

1137 "subdir/001/tiny_problems.dag.metrics", 

1138 "subdir/001/tiny_problems.dag.nodes.log", 

1139 "subdir/001/tiny_problems.node_status", 

1140 }, 

1141 ) 

1142 

1143 def testRelativeSubdir(self): 

1144 with temporaryDirectory() as tmp_dir: 

1145 test_tmp_dir = pathlib.Path(tmp_dir) 

1146 submit_dir = test_tmp_dir / "submit" 

1147 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1148 lssthtc.htc_backup_files(submit_dir, "reldir") 

1149 result_submit = [] 

1150 for root, _, files in os.walk(submit_dir): 

1151 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files]) 

1152 self.assertEqual( 

1153 set(result_submit), 

1154 { 

1155 "./tiny_problems.dag.dagman.log", 

1156 "./tiny_problems.dag.dagman.out", 

1157 "./tiny_problems.dag", 

1158 "./tiny_problems.dag.rescue001", 

1159 "reldir/001/tiny_problems.info.json", 

1160 "reldir/001/tiny_problems.dag.metrics", 

1161 "reldir/001/tiny_problems.dag.nodes.log", 

1162 "reldir/001/tiny_problems.node_status", 

1163 }, 

1164 ) 

1165 

1166 def testSubdags(self): 

1167 with temporaryDirectory() as tmp_dir: 

1168 test_tmp_dir = pathlib.Path(tmp_dir) 

1169 submit_dir = test_tmp_dir / "submit" 

1170 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1171 lssthtc.htc_backup_files(submit_dir) 

1172 result_submit = [] 

1173 for root, _, files in os.walk(submit_dir): 

1174 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files]) 

1175 self.assertEqual( 

1176 set(result_submit), 

1177 { 

1178 "./group_failed_1.dag", 

1179 "./group_failed_1.dag.dagman.log", 

1180 "./group_failed_1.dag.dagman.out", 

1181 "./group_failed_1.dag.rescue001", 

1182 "subdags/wms_group_order1_val1a/group_order1_val1a.dag", 

1183 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.dagman.log", 

1184 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.dagman.out", 

1185 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.nodes.log", 

1186 "subdags/wms_group_order1_val1a/group_order1_val1a.node_status", 

1187 "subdags/wms_group_order1_val1a/wms_group_order1_val1a.dag.post.out", 

1188 "subdags/wms_group_order1_val1a/wms_group_order1_val1a.status.txt", 

1189 "subdags/wms_group_order1_val1b/group_order1_val1b.dag", 

1190 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.dagman.log", 

1191 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.dagman.out", 

1192 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.rescue001", 

1193 "subdags/wms_group_order1_val1c/group_order1_val1c.dag", 

1194 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.dagman.log", 

1195 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.dagman.out", 

1196 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.nodes.log", 

1197 "subdags/wms_group_order1_val1c/group_order1_val1c.node_status", 

1198 "subdags/wms_group_order1_val1c/wms_group_order1_val1c.dag.post.out", 

1199 "subdags/wms_group_order1_val1c/wms_group_order1_val1c.status.txt", 

1200 "001/group_failed_1.dag.nodes.log", 

1201 "001/group_failed_1.info.json", 

1202 "001/group_failed_1.node_status", 

1203 "001/subdags/wms_group_order1_val1b/group_order1_val1b.dag.nodes.log", 

1204 "001/subdags/wms_group_order1_val1b/group_order1_val1b.node_status", 

1205 "001/subdags/wms_group_order1_val1b/wms_group_order1_val1b.status.txt", 

1206 "001/subdags/wms_group_order1_val1b/wms_group_order1_val1b.dag.post.out", 

1207 }, 

1208 ) 

1209 

1210 

1211class UpdateRescueFileTestCase(unittest.TestCase): 

1212 """Test _update_rescue_file function.""" 

1213 

1214 def testSuccess(self): 

1215 self.maxDiff = None 

1216 with temporaryDirectory() as tmp_dir: 

1217 test_tmp_dir = pathlib.Path(tmp_dir) 

1218 submit_dir = test_tmp_dir / "submit" 

1219 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1220 rescue_file = submit_dir / "group_failed_1.dag.rescue001" 

1221 failed_subdags = lssthtc._update_rescue_file(rescue_file) 

1222 self.assertEqual(set(failed_subdags), {"wms_group_order1_val1b"}) 

1223 with open(rescue_file) as fh: 

1224 lines = fh.readlines() 

1225 results = "".join(lines) 

1226 

1227 truth = """# Rescue DAG file, created after running 

1228# the u_testuser_DM-46294_group_fail_20250310T160455Z.dag DAG file 

1229# Created 3/10/2025 16:08:56 UTC 

1230# Rescue DAG version: 2.0.1 (partial) 

1231# 

1232# Total number of Nodes: 26 

1233# Nodes premarked DONE: 21 

1234# Nodes that failed: 2 

1235# wms_group_order1_val1b,finalJob,<ENDLIST> 

1236 

1237DONE pipetaskInit 

1238DONE label1_val1c_val2a 

1239DONE label1_val1b_val2b 

1240DONE label1_val1b_val2a 

1241DONE label1_val1c_val2b 

1242DONE label1_val1a_val2a 

1243DONE label1_val1a_val2b 

1244DONE label3_val1c_val2a 

1245DONE label3_val1b_val2b 

1246DONE label3_val1b_val2a 

1247DONE label3_val1c_val2b 

1248DONE label3_val1a_val2a 

1249DONE label3_val1a_val2b 

1250DONE wms_group_order1_val1a 

1251DONE label5_val1a_val2a 

1252DONE label5_val1a_val2b 

1253DONE wms_group_order1_val1c 

1254DONE label5_val1c_val2a 

1255DONE label5_val1c_val2b 

1256DONE wms_check_status_wms_group_order1_val1a 

1257DONE wms_check_status_wms_group_order1_val1c 

1258""" 

1259 

1260 self.assertEqual(results, truth) 

1261 

1262 

1263class ReadRescueHeadersTestCase(unittest.TestCase): 

1264 """Test _read_rescue_headers function.""" 

1265 

1266 def testTypical(self): 

1267 content = "# Header line 1\n# Header line 2\n\nDONE somenode\n" 

1268 result = lssthtc._read_rescue_headers(io.StringIO(content)) 

1269 self.assertEqual(result, ["# Header line 1", "# Header line 2"]) 

1270 

1271 def testEmptyFile(self): 

1272 result = lssthtc._read_rescue_headers(io.StringIO("")) 

1273 self.assertEqual(result, []) 

1274 

1275 def testOnlyHeaderLines(self): 

1276 content = "# Line 1\n# Line 2\n# Line 3\n" 

1277 result = lssthtc._read_rescue_headers(io.StringIO(content)) 

1278 self.assertEqual(result, ["# Line 1", "# Line 2", "# Line 3"]) 

1279 

1280 def testFirstLineNotComment(self): 

1281 content = "DONE somenode\n# Header\n" 

1282 result = lssthtc._read_rescue_headers(io.StringIO(content)) 

1283 self.assertEqual(result, []) 

1284 

1285 def testWhitespaceStripped(self): 

1286 content = " # Header line 1 \n # Header line 2 \n\n" 

1287 result = lssthtc._read_rescue_headers(io.StringIO(content)) 

1288 self.assertEqual(result, ["# Header line 1", "# Header line 2"]) 

1289 

1290 

1291class WriteRescueHeadersTestCase(unittest.TestCase): 

1292 """Test _write_rescue_headers function.""" 

1293 

1294 def testTypical(self): 

1295 header_lines = ["# Header line 1", "# Header line 2", "# Header line 3"] 

1296 outfh = io.StringIO() 

1297 lssthtc._write_rescue_headers(header_lines, outfh) 

1298 self.assertEqual(outfh.getvalue(), "# Header line 1\n# Header line 2\n# Header line 3\n\n") 

1299 

1300 def testEmptyList(self): 

1301 outfh = io.StringIO() 

1302 lssthtc._write_rescue_headers([], outfh) 

1303 self.assertEqual(outfh.getvalue(), "\n") 

1304 

1305 def testSingleLine(self): 

1306 outfh = io.StringIO() 

1307 lssthtc._write_rescue_headers(["# Only line"], outfh) 

1308 self.assertEqual(outfh.getvalue(), "# Only line\n\n") 

1309 

1310 

1311class UpdateRescueHeadersTestCase(unittest.TestCase): 

1312 """Test _update_rescue_headers function.""" 

1313 

1314 def testWithFailedSubdag(self): 

1315 header_lines = [ 

1316 "# Total number of Nodes: 26", 

1317 "# Nodes premarked DONE: 22", 

1318 "# Nodes that failed: 2", 

1319 "# wms_check_status_wms_group_order1_val1b,finalJob,<ENDLIST>", 

1320 ] 

1321 result = lssthtc._update_rescue_headers(header_lines) 

1322 self.assertEqual(result, ["wms_group_order1_val1b"]) 

1323 self.assertEqual(header_lines[1], "# Nodes premarked DONE: 21") 

1324 self.assertEqual(header_lines[3], "# wms_group_order1_val1b,finalJob,<ENDLIST>") 

1325 

1326 def testNoSubdagFailures(self): 

1327 header_lines = [ 

1328 "# Nodes premarked DONE: 5", 

1329 "# Nodes that failed: 1", 

1330 "# finalJob,<ENDLIST>", 

1331 ] 

1332 result = lssthtc._update_rescue_headers(header_lines) 

1333 self.assertEqual(result, []) 

1334 self.assertEqual(header_lines[0], "# Nodes premarked DONE: 5") 

1335 self.assertEqual(header_lines[2], "# finalJob,<ENDLIST>") 

1336 

1337 def testMultipleFailedSubdags(self): 

1338 header_lines = [ 

1339 "# Nodes premarked DONE: 10", 

1340 "# Nodes that failed: 3", 

1341 "# wms_check_status_subdag_a,wms_check_status_subdag_b,finalJob,<ENDLIST>", 

1342 ] 

1343 result = lssthtc._update_rescue_headers(header_lines) 

1344 self.assertEqual(result, ["subdag_a", "subdag_b"]) 

1345 self.assertEqual(header_lines[0], "# Nodes premarked DONE: 8") 

1346 self.assertEqual(header_lines[2], "# subdag_a,subdag_b,finalJob,<ENDLIST>") 

1347 

1348 def testNoFailedNodesLine(self): 

1349 header_lines = [ 

1350 "# Total number of Nodes: 5", 

1351 "# Nodes premarked DONE: 5", 

1352 ] 

1353 original = list(header_lines) 

1354 result = lssthtc._update_rescue_headers(header_lines) 

1355 self.assertEqual(result, []) 

1356 self.assertEqual(header_lines, original) 

1357 

1358 def testEmptyHeader(self): 

1359 result = lssthtc._update_rescue_headers([]) 

1360 self.assertEqual(result, []) 

1361 

1362 

1363class ReadDagStatusTestCase(unittest.TestCase): 

1364 """Test read_dag_status function and read_single_dag_status.""" 

1365 

1366 def testFileMissing(self): 

1367 with temporaryDirectory() as tmp_dir: 

1368 with self.assertRaisesRegex(FileNotFoundError, "DAGMan node status not found"): 

1369 _ = lssthtc.read_dag_status(tmp_dir) 

1370 

1371 def testRegular(self): 

1372 with temporaryDirectory() as tmp_dir: 

1373 submit_dir = os.path.join(tmp_dir, "tiny_problems") 

1374 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1375 results = lssthtc.read_dag_status(submit_dir) 

1376 truth = { 

1377 "JobProcsHeld": 0, 

1378 "NodesPost": 0, 

1379 "JobProcsIdle": 0, 

1380 "NodesTotal": 6, 

1381 "NodesFailed": 2, 

1382 "NodesDone": 3, 

1383 "NodesQueued": 0, 

1384 "NodesPre": 0, 

1385 "NodesFutile": 1, 

1386 "NodesUnready": 0, 

1387 } 

1388 self.assertEqual(results, results | truth) 

1389 

1390 def testSubdags(self): 

1391 """Making sure it gets data from subdag dirs and doesn't 

1392 fail if some subdags haven't started running yet. 

1393 """ 

1394 self.maxDiff = None 

1395 with temporaryDirectory() as tmp_dir: 

1396 submit_dir = os.path.join(tmp_dir, "submit") 

1397 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*")) 

1398 results = lssthtc.read_dag_status(submit_dir) 

1399 truth = { 

1400 "JobProcsHeld": 0, 

1401 "NodesPost": 0, 

1402 "JobProcsIdle": 0, 

1403 "NodesTotal": 34, 

1404 "NodesFailed": 0, 

1405 "NodesDone": 17, 

1406 "NodesQueued": 3, 

1407 "NodesPre": 0, 

1408 "NodesFutile": 0, 

1409 "NodesUnready": 14, 

1410 } 

1411 self.assertEqual(results, results | truth) 

1412 

1413 

1414class ReadDagInfoTestCase(unittest.TestCase): 

1415 """Test read_dag_info function.""" 

1416 

1417 def testFileMissing(self): 

1418 with temporaryDirectory() as tmp_dir: 

1419 with self.assertRaisesRegex(FileNotFoundError, "File with DAGMan job information not found in "): 

1420 _ = lssthtc.read_dag_info(tmp_dir) 

1421 

1422 def testSuccess(self): 

1423 with temporaryDirectory() as tmp_dir: 

1424 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.info.json", tmp_dir) 

1425 filename, results = lssthtc.read_dag_info(tmp_dir) 

1426 self.assertIn("info.json", str(filename)) 

1427 self.assertTrue(pathlib.Path(filename).is_file()) 

1428 

1429 truth = { 

1430 "test02": { 

1431 "9208.0": { 

1432 "ClusterId": 9208, 

1433 "GlobalJobId": "test02#9208.0#1739465078", 

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

1435 "bps_project": "dev", 

1436 "bps_payload": "tiny", 

1437 "bps_operator": "testuser", 

1438 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorWorkflow", 

1439 "bps_provisioning_job": "provisioningJob", 

1440 "bps_run_quanta": "label1:1;label2:1", 

1441 "bps_campaign": "quick", 

1442 "bps_runsite": "testpool", 

1443 "bps_job_summary": "pipetaskInit:1;label1:1;label2:1;finalJob:1", 

1444 "bps_run": "u_testuser_tiny_20250213T164427Z", 

1445 "bps_isjob": "True", 

1446 } 

1447 } 

1448 } 

1449 

1450 self.assertEqual(results, truth) 

1451 

1452 def testPermissionError(self): 

1453 with temporaryDirectory() as tmp_dir: 

1454 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.info.json", tmp_dir) 

1455 with unittest.mock.patch("lsst.ctrl.bps.htcondor.lssthtc.open") as mocked_open: 

1456 mocked_open.side_effect = PermissionError 

1457 with self.assertLogs("lsst.ctrl.bps.htcondor", level="DEBUG") as cm: 

1458 _, results = lssthtc.read_dag_info(tmp_dir) 

1459 self.assertIn("Retrieving DAGMan job information failed:", cm.output[-1]) 

1460 self.assertEqual({}, results) 

1461 

1462 

1463class HtcWriteCondorFileTestCase(unittest.TestCase): 

1464 """Test htc_write_condor_file function.""" 

1465 

1466 def testSuccess(self): 

1467 with temporaryDirectory() as tmp_dir: 

1468 job_name = "job1" 

1469 filename = pathlib.Path(tmp_dir) / f"label1/{job_name}.sub" 

1470 job = { 

1471 "executable": "$(CTRL_MPEXEC_DIR)/bin/pipetask", 

1472 "arguments": "-a -b 2 -c", 

1473 "request_memory": "2000", 

1474 "environment": "one=1 two=\"2\" three='spacey 'quoted' value'", 

1475 "log": f"{job_name}.log", 

1476 } 

1477 job_attrs = { 

1478 "bps_job_name": job_name, 

1479 "bps_job_label": "label1", 

1480 "bps_job_quanta": "task1:8;task2:8", 

1481 } 

1482 expected = [ 

1483 "executable=$(CTRL_MPEXEC_DIR)/bin/pipetask\n", 

1484 'arguments="-a -b 2 -c"\n', 

1485 "request_memory=2000\n", 

1486 "environment=\"one=1 two=\"2\" three='spacey 'quoted' value'\"\n", 

1487 f"output={job_name}.$(Cluster).out\n", 

1488 f"error={job_name}.$(Cluster).out\n", 

1489 f"log={job_name}.log\n", 

1490 f'+bps_job_name = "{job_name}"\n', 

1491 '+bps_job_label = "label1"\n', 

1492 '+bps_job_quanta = "task1:8;task2:8"\n', 

1493 "queue\n", 

1494 ] 

1495 

1496 lssthtc.htc_write_condor_file(filename, job_name, job, job_attrs) 

1497 with open(filename, encoding="utf-8") as f: 

1498 actual = f.readlines() 

1499 

1500 self.assertEqual(set(actual), set(expected)) 

1501 self.assertTrue(filename.exists()) 

1502 # Try to make Submit object from file to find any syntax issues 

1503 _ = lssthtc.htc_create_submit_from_file(filename) 

1504 

1505 

1506class HtcCreateSubmitFromDagTestCase(unittest.TestCase): 

1507 """Test htc_create_submit_from_dag function.""" 

1508 

1509 @classmethod 

1510 def setUpClass(cls): 

1511 cls.bindir = None 

1512 # htcondor.Submit.from_dag requires condor_dagman executable in path. 

1513 if not which("condor_dagman"): # pragma: no cover 

1514 cls.bindir = tempfile.TemporaryDirectory() 

1515 fake_dagman_exec = pathlib.Path(cls.bindir.name) / "condor_dagman" 

1516 with open(fake_dagman_exec, "w") as fh: 

1517 print("#!/bin/bash", file=fh) 

1518 print("echo fake_condor_dagman $@", file=fh) 

1519 print("exit 0", file=fh) 

1520 fake_dagman_exec.chmod(fake_dagman_exec.stat().st_mode | stat.S_IEXEC) 

1521 os.environ["PATH"] = f"{os.environ['PATH']}:{cls.bindir.name}" 

1522 

1523 @classmethod 

1524 def tearDownClass(cls): 

1525 if cls.bindir: 1525 ↛ 1526line 1525 didn't jump to line 1526 because the condition on line 1525 was never true

1526 cls.bindir.cleanup() 

1527 

1528 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"}) 

1529 def testMaxIdleEnvVar(self): 

1530 with temporaryDirectory() as tmp_dir: 

1531 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir) 

1532 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag" 

1533 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}) 

1534 self.assertIn("-MaxIdle 42", submit["arguments"]) 

1535 

1536 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"}) 

1537 def testMaxIdleInDAGManConfig(self): 

1538 with temporaryDirectory() as tmp_dir: 

1539 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir) 

1540 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag" 

1541 config_filename = pathlib.Path(tmp_dir) / "dagman.conf" 

1542 with open(config_filename, "w") as fh: 

1543 print("DAGMAN_MAX_JOBS_IDLE = 300", file=fh) 

1544 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, config_filename) 

1545 self.assertIn("-MaxIdle 300", submit["arguments"]) 

1546 

1547 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"}) 

1548 def testMaxIdleNotInDAGManConfig(self): 

1549 with temporaryDirectory() as tmp_dir: 

1550 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir) 

1551 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag" 

1552 config_filename = pathlib.Path(tmp_dir) / "dagman.conf" 

1553 with open(config_filename, "w") as fh: 

1554 print("DAGMAN_MAX_JOBS_SUBMITTED = 300", file=fh) 

1555 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, config_filename) 

1556 self.assertIn("-MaxIdle 42", submit["arguments"]) 

1557 

1558 @unittest.mock.patch.dict(os.environ, {}) 

1559 def testMaxIdleGiven(self): 

1560 with temporaryDirectory() as tmp_dir: 

1561 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir) 

1562 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag" 

1563 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {"MaxIdle": 37}) 

1564 self.assertIn("-MaxIdle 37", submit["arguments"]) 

1565 

1566 @unittest.mock.patch.dict(os.environ, {}) 

1567 def testMaxJobsIdleParam(self): 

1568 def _fake_params_contains(key): 

1569 if key == "DAGMAN_MAX_JOBS_IDLE": 

1570 return True 

1571 return False # pragma: no cover 

1572 

1573 def _fake_params_get(key): 

1574 if key == "DAGMAN_MAX_JOBS_IDLE": 

1575 return 16 

1576 return "FAKE_VAL" # pragma: no cover 

1577 

1578 with temporaryDirectory() as tmp_dir: 

1579 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir) 

1580 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag" 

1581 with unittest.mock.patch("htcondor.param") as mock_param: 

1582 mock_param.__contains__.side_effect = _fake_params_contains 

1583 mock_param.__getitem__.side_effect = _fake_params_get 

1584 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, None) 

1585 self.assertIn("-MaxIdle 16", submit["arguments"]) 

1586 

1587 @unittest.mock.patch.dict(os.environ, {}) 

1588 def testNoMaxJobsIdle(self): 

1589 """Note: Since the produced arguments differ depending on 

1590 HTCondor version when no MaxIdle passed to from_dag, not 

1591 checking arguments string here. Instead just making sure 

1592 lssthtc code doesn't pass MaxIdle value to from_dag. 

1593 """ 

1594 with temporaryDirectory() as tmp_dir: 

1595 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir) 

1596 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag" 

1597 with unittest.mock.patch("htcondor.Submit.from_dag") as submit_mock: 

1598 with unittest.mock.patch("htcondor.param") as mock_param: 

1599 mock_param.__contains__.return_value = False 

1600 _ = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}) 

1601 submit_mock.assert_called_once_with(str(dag_filename), {}) 

1602 

1603 

1604class HtcDagTestCase(unittest.TestCase): 

1605 """Test for HTCDag class.""" 

1606 

1607 def setUp(self): 

1608 job = lssthtc.HTCJob(name="test_job") 

1609 job.add_job_cmds( 

1610 { 

1611 "executable": "/usr/bin/echo", 

1612 "arguments": "foo", 

1613 "output": "test_job.$(Cluster).out", 

1614 "error": "test_job.$(Cluster).out", 

1615 "log": "test_job.$(Cluster).log", 

1616 } 

1617 ) 

1618 job.subfile = f"{job.name}.sub" 

1619 

1620 self.dag = lssthtc.HTCDag(name="test_workflow") 

1621 self.dag.add_job(job) 

1622 

1623 self.subfile_expected = [ 

1624 "executable=/usr/bin/echo\n", 

1625 'arguments="foo"\n', 

1626 "output=test_job.$(Cluster).out\n", 

1627 "error=test_job.$(Cluster).out\n", 

1628 "log=test_job.$(Cluster).log\n", 

1629 "queue\n", 

1630 ] 

1631 

1632 def tearDown(self): 

1633 pass 

1634 

1635 def testWriteWithDagConfig(self): 

1636 with temporaryDirectory() as tmp_dir: 

1637 config = BpsConfig(Config(htcondor_config.HTC_DEFAULTS_URI)) 

1638 job = self.dag.nodes["test_job"]["data"] 

1639 wms_config_filename = "dagman.conf" 

1640 wms_configurator = dagman_configurator.DagmanConfigurator(config) 

1641 wms_configurator.prepare(wms_config_filename, prefix=tmp_dir) 

1642 wms_configurator.configure(self.dag) 

1643 dagfile_expected = [ 

1644 f"CONFIG {wms_config_filename}\n", 

1645 f'JOB {job.name} "{job.subfile}"\n', 

1646 f"NODE_STATUS_FILE {self.dag.name}.node_status\n", 

1647 f'SET_JOB_ATTR bps_wms_config_path= "{wms_config_filename}"\n', 

1648 ] 

1649 

1650 self.dag.write(tmp_dir, "", "") 

1651 

1652 self.assertIn("submit_path", self.dag.graph) 

1653 self.assertEqual(self.dag.graph["submit_path"], tmp_dir) 

1654 self.assertIn("dag_filename", self.dag.graph) 

1655 self.assertEqual(self.dag.graph["dag_filename"], f"{self.dag.graph['name']}.dag") 

1656 with open(os.path.join(tmp_dir, self.dag.graph["dag_filename"]), encoding="utf-8") as f: 

1657 dagfile_actual = f.readlines() 

1658 self.assertEqual(dagfile_actual, dagfile_expected) 

1659 with open(os.path.join(tmp_dir, job.subfile), encoding="utf-8") as f: 

1660 subfile_actual = f.readlines() 

1661 self.assertEqual(subfile_actual, self.subfile_expected) 

1662 

1663 def testWriteWithoutDagConfig(self): 

1664 with temporaryDirectory() as tmp_dir: 

1665 job = self.dag.nodes["test_job"]["data"] 

1666 dagfile_expected = [ 

1667 f'JOB {job.name} "{job.subfile}"\n', 

1668 f"NODE_STATUS_FILE {self.dag.name}.node_status\n", 

1669 ] 

1670 

1671 self.dag.write(tmp_dir, "", "") 

1672 

1673 self.assertIn("submit_path", self.dag.graph) 

1674 self.assertEqual(self.dag.graph["submit_path"], tmp_dir) 

1675 self.assertIn("dag_filename", self.dag.graph) 

1676 self.assertEqual(self.dag.graph["dag_filename"], f"{self.dag.graph['name']}.dag") 

1677 with open(os.path.join(tmp_dir, self.dag.graph["dag_filename"]), encoding="utf-8") as f: 

1678 dagfile_actual = f.readlines() 

1679 self.assertEqual(dagfile_actual, dagfile_expected) 

1680 with open(os.path.join(tmp_dir, job.subfile), encoding="utf-8") as f: 

1681 subfile_actual = f.readlines() 

1682 self.assertEqual(subfile_actual, self.subfile_expected) 

1683 

1684 def testWriteLazySubdag(self): 

1685 self.maxDiff = None 

1686 dag, truth_files = make_lazy_dag("test1", True) 

1687 dag.graph["write_dot"] = True 

1688 with temporaryDirectory() as tmp_dir: 

1689 dag.write(tmp_dir, "", "") 

1690 with open(os.path.join(tmp_dir, dag.graph["dag_filename"]), encoding="utf-8") as f: 

1691 dagfile_actual = f.readlines() 

1692 self.assertIn("DOT test1.dot\n", dagfile_actual) 

1693 

1694 all_files = [] 

1695 for root, _, files in pathlib.Path(tmp_dir).walk(): 

1696 relroot = pathlib.Path(root).relative_to(tmp_dir) 

1697 all_files.extend([str(relroot / f) for f in files]) 

1698 self.assertEqual(sorted(all_files), sorted(truth_files)) 

1699 

1700 @staticmethod 

1701 def _make_simple_job(name): 

1702 job = lssthtc.HTCJob(name=name) 

1703 job.add_job_cmds({"executable": "/usr/bin/echo", "arguments": name}) 

1704 job.subfile = f"{name}.sub" 

1705 return job 

1706 

1707 def testWriteEdgesAndSpecialJobs(self): 

1708 self.maxDiff = None 

1709 dag = lssthtc.HTCDag(name="test_special") 

1710 dag.add_attribs({"bps_run": "myrun"}) 

1711 dag.add_job(self._make_simple_job("jobA")) 

1712 dag.add_job(self._make_simple_job("jobB")) 

1713 dag.add_job_relationships(["jobA"], ["jobB"]) 

1714 dag.add_final_job(self._make_simple_job("finalJob")) 

1715 dag.add_service_job(self._make_simple_job("svc")) 

1716 

1717 dagfile_expected = [ 

1718 'JOB jobA "jobA.sub"\n', 

1719 'JOB jobB "jobB.sub"\n', 

1720 "PARENT jobA CHILD jobB\n", 

1721 f"NODE_STATUS_FILE {dag.name}.node_status\n", 

1722 'SET_JOB_ATTR bps_run= "myrun"\n', 

1723 'FINAL finalJob "finalJob.sub"\n', 

1724 'SERVICE svc "svc.sub"\n', 

1725 ] 

1726 

1727 with temporaryDirectory() as tmp_dir: 

1728 dag.write(tmp_dir, "", "") 

1729 with open(os.path.join(tmp_dir, dag.graph["dag_filename"]), encoding="utf-8") as f: 

1730 self.assertEqual(f.readlines(), dagfile_expected) 

1731 # Submit files are written for regular and special jobs. 

1732 for subfile in ("jobA.sub", "jobB.sub", "finalJob.sub", "svc.sub"): 

1733 self.assertTrue(os.path.exists(os.path.join(tmp_dir, subfile))) 

1734 

1735 def testWriteDagSubdir(self): 

1736 # A non-empty dag_subdir places the .dag file (and DIR clause) under 

1737 # that subdirectory, creating it as needed. 

1738 dag = lssthtc.HTCDag(name="test_subdir") 

1739 dag.add_job(self._make_simple_job("jobA")) 

1740 

1741 with temporaryDirectory() as tmp_dir: 

1742 dag.write(tmp_dir, "", "daglevel") 

1743 self.assertEqual(dag.graph["dag_filename"], "daglevel/test_subdir.dag") 

1744 self.assertTrue(os.path.exists(os.path.join(tmp_dir, "daglevel", "test_subdir.dag"))) 

1745 

1746 def testWriteMissingJobData(self): 

1747 # A node without the "data" key should raise KeyError. 

1748 dag = lssthtc.HTCDag(name="test_baddata") 

1749 dag.add_node("orphan") # no data=... provided 

1750 with temporaryDirectory() as tmp_dir: 

1751 with self.assertRaises(KeyError): 

1752 dag.write(tmp_dir, "", "") 

1753 

1754 

1755class WriteDagInfoTestCase(unittest.TestCase): 

1756 """Test for write_dag_info function.""" 

1757 

1758 def setUp(self): 

1759 self.run = "u_testuser_DM-53494_20260220T001651Z" 

1760 self.data = { 

1761 "mycomputer": { 

1762 "24390.0": { 

1763 "ClusterId": 24390, 

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

1765 "bps_run": self.run, 

1766 "bps_isjob": "True", 

1767 "bps_payload": "DM-53494", 

1768 "bps_project": "dev", 

1769 "bps_runsite": "site1", 

1770 "bps_campaign": "ci_rc2", 

1771 "bps_operator": "testuser", 

1772 "bps_run_quanta": "", 

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

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

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

1776 "bps_wms_config_path": "dagman.conf", 

1777 } 

1778 } 

1779 } 

1780 

1781 def testWrite(self): 

1782 with temporaryDirectory() as tmp_dir: 

1783 with chdir(tmp_dir): 

1784 path = pathlib.Path(tmp_dir) / "test.info.json" 

1785 filename = lssthtc.write_dag_info(path, self.data) 

1786 self.assertTrue(path.is_file(), f"File not found at {path}") 

1787 self.assertEqual(filename, path) 

1788 

1789 read_filename, read_data = lssthtc.read_dag_info(tmp_dir) 

1790 self.assertEqual(read_filename, path) 

1791 self.assertEqual(read_data, self.data) 

1792 

1793 

1794if __name__ == "__main__": 

1795 unittest.main()