Coverage for tests/test_lssthtc.py: 99%

852 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-09-17 02:08 -0700

1# This file is part of ctrl_bps_htcondor. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

5# (https://www.lsst.org). 

6# See the COPYRIGHT file at the top-level directory of this distribution 

7# for details of code ownership. 

8# 

9# This software is dual licensed under the GNU General Public License and also 

10# under a 3-clause BSD license. Recipients may choose which of these licenses 

11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt, 

12# respectively. If you choose the GPL option then the following text applies 

13# (but note that there is still no warranty even if you opt for BSD instead): 

14# 

15# This program is free software: you can redistribute it and/or modify 

16# it under the terms of the GNU General Public License as published by 

17# the Free Software Foundation, either version 3 of the License, or 

18# (at your option) any later version. 

19# 

20# This program is distributed in the hope that it will be useful, 

21# but WITHOUT ANY WARRANTY; without even the implied warranty of 

22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the 

23# GNU General Public License for more details. 

24# 

25# You should have received a copy of the GNU General Public License 

26# along with this program. If not, see <https://www.gnu.org/licenses/>. 

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

225 with temporaryDirectory() as tmp_dir: 

226 with open(os.path.join(tmp_dir, "test_missing_warning.dag.dagman.out"), "w") as fh: 

227 print( 

228 "01/01/26 01:01:01 ERROR: Warning is fatal error because of DAGMAN_USE_STRICT setting", 

229 file=fh, 

230 ) 

231 results = lssthtc.htc_check_dagman_output(tmp_dir) 

232 self.assertIn("Missing warning", results) 

233 

234 

235class SummarizeDagTestCase(unittest.TestCase): 

236 """Test summarize_dag function.""" 

237 

238 def test_no_dag_file(self): 

239 with temporaryDirectory() as tmp_dir: 

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

241 self.assertFalse(len(job_name_to_pipetask)) 

242 self.assertFalse(len(job_name_to_type)) 

243 self.assertFalse(summary) 

244 

245 def test_success(self): 

246 with temporaryDirectory() as tmp_dir: 

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

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

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

250 self.assertEqual( 

251 job_name_to_label, 

252 { 

253 "pipetaskInit": "pipetaskInit", 

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

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

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

257 "finalJob": "finalJob", 

258 }, 

259 ) 

260 self.assertEqual( 

261 job_name_to_type, 

262 { 

263 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

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

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

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

267 "finalJob": lssthtc.WmsNodeType.FINAL, 

268 }, 

269 ) 

270 

271 def test_service(self): 

272 with temporaryDirectory() as tmp_dir: 

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

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

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

276 self.assertEqual( 

277 job_name_to_label, 

278 { 

279 "pipetaskInit": "pipetaskInit", 

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

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

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

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

284 "finalJob": "finalJob", 

285 "provisioningJob": "provisioningJob", 

286 }, 

287 ) 

288 self.assertEqual( 

289 job_name_to_type, 

290 { 

291 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

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

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

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

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

296 "finalJob": lssthtc.WmsNodeType.FINAL, 

297 "provisioningJob": lssthtc.WmsNodeType.SERVICE, 

298 }, 

299 ) 

300 

301 def test_noop(self): 

302 with temporaryDirectory() as tmp_dir: 

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

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

305 self.assertEqual( 

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

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

308 ) 

309 self.assertEqual( 

310 job_name_to_label, 

311 { 

312 "label1_val1a_val2a": "label1", 

313 "label1_val1a_val2b": "label1", 

314 "label1_val1b_val2a": "label1", 

315 "label1_val1b_val2b": "label1", 

316 "label1_val1c_val2a": "label1", 

317 "label1_val1c_val2b": "label1", 

318 "label2_val1a_val2a": "label2", 

319 "label2_val1a_val2b": "label2", 

320 "label2_val1b_val2a": "label2", 

321 "label2_val1b_val2b": "label2", 

322 "label2_val1c_val2a": "label2", 

323 "label2_val1c_val2b": "label2", 

324 "label3_val1a_val2a": "label3", 

325 "label3_val1a_val2b": "label3", 

326 "label3_val1b_val2a": "label3", 

327 "label3_val1b_val2b": "label3", 

328 "label3_val1c_val2a": "label3", 

329 "label3_val1c_val2b": "label3", 

330 "label4_val1a_val2a": "label4", 

331 "label4_val1a_val2b": "label4", 

332 "label4_val1b_val2a": "label4", 

333 "label4_val1b_val2b": "label4", 

334 "label4_val1c_val2a": "label4", 

335 "label4_val1c_val2b": "label4", 

336 "label5_val1a_val2a": "label5", 

337 "label5_val1a_val2b": "label5", 

338 "label5_val1b_val2a": "label5", 

339 "label5_val1b_val2b": "label5", 

340 "label5_val1c_val2a": "label5", 

341 "label5_val1c_val2b": "label5", 

342 "finalJob": "finalJob", 

343 "pipetaskInit": "pipetaskInit", 

344 "wms_noop_order1_val1a": "order1", 

345 "wms_noop_order1_val1b": "order1", 

346 }, 

347 ) 

348 self.assertEqual( 

349 job_name_to_type, 

350 { 

351 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

352 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

353 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

354 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

355 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

356 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

357 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

358 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

359 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

360 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

361 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

362 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

363 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

364 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

365 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

366 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

367 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

368 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

369 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

370 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

371 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

372 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

373 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

374 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

375 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

376 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

377 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

378 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

379 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

380 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

381 "finalJob": lssthtc.WmsNodeType.FINAL, 

382 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

383 "wms_noop_order1_val1a": lssthtc.WmsNodeType.NOOP, 

384 "wms_noop_order1_val1b": lssthtc.WmsNodeType.NOOP, 

385 }, 

386 ) 

387 

388 def test_subdags(self): 

389 self.maxDiff = None 

390 with temporaryDirectory() as tmp_dir: 

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

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

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

394 self.assertEqual( 

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

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

397 ) 

398 

399 self.assertEqual( 

400 job_name_to_label, 

401 { 

402 "pipetaskInit": "pipetaskInit", 

403 "label1_val1b_val2a": "label1", 

404 "label1_val1c_val2a": "label1", 

405 "label1_val1a_val2b": "label1", 

406 "label1_val1b_val2b": "label1", 

407 "label1_val1c_val2b": "label1", 

408 "label1_val1a_val2a": "label1", 

409 "label2_val1a_val2b": "label2", 

410 "label2_val1a_val2a": "label2", 

411 "label2_val1b_val2a": "label2", 

412 "label2_val1b_val2b": "label2", 

413 "label2_val1c_val2a": "label2", 

414 "label2_val1c_val2b": "label2", 

415 "label3_val1b_val2a": "label3", 

416 "label3_val1c_val2a": "label3", 

417 "label3_val1a_val2b": "label3", 

418 "label3_val1b_val2b": "label3", 

419 "label3_val1c_val2b": "label3", 

420 "label3_val1a_val2a": "label3", 

421 "label4_val1a_val2b": "label4", 

422 "label4_val1a_val2a": "label4", 

423 "label4_val1b_val2a": "label4", 

424 "label4_val1b_val2b": "label4", 

425 "label4_val1c_val2a": "label4", 

426 "label4_val1c_val2b": "label4", 

427 "label5_val1a_val2b": "label5", 

428 "label5_val1a_val2a": "label5", 

429 "label5_val1b_val2a": "label5", 

430 "label5_val1b_val2b": "label5", 

431 "label5_val1c_val2a": "label5", 

432 "label5_val1c_val2b": "label5", 

433 "finalJob": "finalJob", 

434 "provisioningJob": "provisioningJob", 

435 "wms_group_order1_val1a": "order1", 

436 "wms_group_order1_val1b": "order1", 

437 "wms_group_order1_val1c": "order1", 

438 "wms_check_status_wms_group_order1_val1a": "order1", 

439 "wms_check_status_wms_group_order1_val1b": "order1", 

440 "wms_check_status_wms_group_order1_val1c": "order1", 

441 }, 

442 ) 

443 

444 self.assertEqual( 

445 job_name_to_type, 

446 { 

447 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

448 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

449 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

450 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

451 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

452 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

453 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

454 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

455 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

456 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

457 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

458 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

459 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

460 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

461 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

462 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

463 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

464 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

465 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

466 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

467 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

468 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

469 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

470 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

471 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

472 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

473 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

474 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

475 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

476 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

477 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

478 "finalJob": lssthtc.WmsNodeType.FINAL, 

479 "provisioningJob": lssthtc.WmsNodeType.SERVICE, 

480 "wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG, 

481 "wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG, 

482 "wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG, 

483 "wms_check_status_wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG_CHECK, 

484 "wms_check_status_wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG_CHECK, 

485 "wms_check_status_wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG_CHECK, 

486 }, 

487 ) 

488 

489 

490class ReadDagNodesLogTestCase(unittest.TestCase): 

491 """Test read_dag_nodes_log function.""" 

492 

493 def setUp(self): 

494 self.tmpdir = tempfile.mkdtemp() 

495 

496 def tearDown(self): 

497 rmtree(self.tmpdir, ignore_errors=True) 

498 

499 def testFileMissing(self): 

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

501 _ = lssthtc.read_dag_nodes_log(self.tmpdir) 

502 

503 def testRegular(self): 

504 with temporaryDirectory() as tmp_dir: 

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

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

507 results = lssthtc.read_dag_nodes_log(submit_dir) 

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

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

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

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

512 

513 def testSubdags(self): 

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

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

516 """ 

517 with temporaryDirectory() as tmp_dir: 

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

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

520 results = lssthtc.read_dag_nodes_log(submit_dir) 

521 # main dag 

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

523 # subdag 

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

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

526 

527 

528class ReadNodeStatusTestCase(unittest.TestCase): 

529 """Test read_node_status function.""" 

530 

531 def setUp(self): 

532 self.tmpdir = tempfile.mkdtemp() 

533 

534 def tearDown(self): 

535 rmtree(self.tmpdir, ignore_errors=True) 

536 

537 def testServiceJobNotSubmitted(self): 

538 # tiny_prov_no_submit files have successful workflow 

539 # but provisioningJob could not submit. 

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

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

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

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

544 

545 jobs = lssthtc.read_node_status(self.tmpdir) 

546 found = [ 

547 id_ 

548 for id_ in jobs 

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

550 ] 

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

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

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

554 

555 def testMissingStatusFile(self): 

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

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

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

559 

560 jobs = lssthtc.read_node_status(self.tmpdir) 

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

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

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

564 found = [ 

565 id_ 

566 for id_ in jobs 

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

568 ] 

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

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

571 

572 def testSubdagsRunning(self): 

573 with temporaryDirectory() as tmp_dir: 

574 test_tmp_dir = pathlib.Path(tmp_dir) 

575 submit_dir = test_tmp_dir / "submit" 

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

577 jobs = lssthtc.read_node_status(submit_dir) 

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

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

580 job_name_to_id = {} 

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

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

583 job_type_to_names = {} 

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

585 job_type_to_names.setdefault( 

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

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

588 

589 # check counts 

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

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

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

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

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

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

596 

597 # spot check some statuses 

598 self.assertEqual( 

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

600 ) 

601 self.assertEqual( 

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

603 ) 

604 self.assertEqual( 

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

606 ) 

607 self.assertEqual( 

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

609 ) 

610 

611 def testSubdagsFailed(self): 

612 with temporaryDirectory() as tmp_dir: 

613 test_tmp_dir = pathlib.Path(tmp_dir) 

614 submit_dir = test_tmp_dir / "submit" 

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

616 jobs = lssthtc.read_node_status(submit_dir) 

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

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

619 job_name_to_id = {} 

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

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

622 job_type_to_names = {} 

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

624 job_type_to_names.setdefault( 

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

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

627 

628 # check counts 

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

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

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

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

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

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

635 

636 # spot check some statuses 

637 self.assertEqual( 

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

639 ) 

640 self.assertEqual( 

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

642 ) 

643 self.assertEqual( 

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

645 ) 

646 

647 self.assertEqual( 

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

649 ) 

650 self.assertEqual( 

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

652 ) 

653 self.assertEqual( 

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

655 lssthtc.NodeStatus.ERROR, 

656 ) 

657 

658 

659class ReadSingleNodeStatusTestCase(unittest.TestCase): 

660 """Test read_single_node_status function.""" 

661 

662 def setUp(self): 

663 self.tmpdir = tempfile.mkdtemp() 

664 

665 def tearDown(self): 

666 rmtree(self.tmpdir, ignore_errors=True) 

667 

668 def _copyFiles(self, data_subdir, suffixes): 

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

670 for suffix in suffixes: 

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

672 

673 def _jobNameToId(self, jobs): 

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

675 

676 def testAllDone(self): 

677 self._copyFiles( 

678 "tiny_success", 

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

680 ) 

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

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

683 

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

685 name_to_id = self._jobNameToId(jobs) 

686 

687 # All four submitted nodes are marked DONE. 

688 for name in [ 

689 "pipetaskInit", 

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

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

692 "finalJob", 

693 ]: 

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

695 self.assertEqual( 

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

697 lssthtc.NodeStatus.DONE, 

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

699 ) 

700 

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

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

703 self.assertIn("provisioningJob", name_to_id) 

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

705 

706 # Spot-check labels and types. 

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

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

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

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

711 

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

713 for job in jobs.values(): 

714 self.assertIn("DAGManJobID", job) 

715 

716 def testMixedStatuses(self): 

717 self._copyFiles( 

718 "tiny_problems", 

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

720 ) 

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

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

723 

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

725 name_to_id = self._jobNameToId(jobs) 

726 

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

728 self.assertEqual( 

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

730 lssthtc.NodeStatus.ERROR, 

731 ) 

732 self.assertEqual( 

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

734 lssthtc.NodeStatus.DONE, 

735 ) 

736 self.assertEqual( 

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

738 lssthtc.NodeStatus.FUTILE, 

739 ) 

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

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

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

743 self.assertIn("provisioningJob", name_to_id) 

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

745 

746 def testRunningWorkflow(self): 

747 self._copyFiles( 

748 "tiny_running", 

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

750 ) 

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

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

753 

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

755 name_to_id = self._jobNameToId(jobs) 

756 

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

758 self.assertEqual( 

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

760 lssthtc.NodeStatus.SUBMITTED, 

761 ) 

762 self.assertEqual( 

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

764 lssthtc.NodeStatus.NOT_READY, 

765 ) 

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

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

768 self.assertIn("provisioningJob", name_to_id) 

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

770 

771 def testMissingNodeStatusFile(self): 

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

773 # and dag. 

774 self._copyFiles( 

775 "tiny_problems", 

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

777 ) 

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

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

780 

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

782 name_to_id = self._jobNameToId(jobs) 

783 

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

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

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

787 

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

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

790 

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

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

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

794 

795 def testMissingLogFiles(self): 

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

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

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

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

800 

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

802 for job in jobs.values(): 

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

804 self.assertEqual(job["DAGManJobID"], lssthtc.MISSING_ID) 

805 

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

807 name_to_id = self._jobNameToId(jobs) 

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

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

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

811 

812 def testPermissionLogFiles(self): 

813 # Check what happens if can't read dag log file. 

814 self._copyFiles("tiny_success", [".dag", ".node_status", ".dag.dagman.log"]) 

815 filename = os.path.join(self.tmpdir, "tiny_success.dag.dagman.log") 

816 current_mode = os.stat(filename).st_mode 

817 no_read_mode = current_mode & ~stat.S_IRUSR & ~stat.S_IRGRP & ~stat.S_IROTH 

818 os.chmod(filename, no_read_mode) 

819 

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

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

822 

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

824 for job in jobs.values(): 

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

826 self.assertEqual(job["DAGManJobID"], lssthtc.MISSING_ID) 

827 

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

829 name_to_id = self._jobNameToId(jobs) 

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

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

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

833 

834 def testInitFakeId(self): 

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

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

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

838 init_fake_id = -10 

839 jobs = lssthtc.read_single_node_status(filename, init_fake_id) 

840 

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

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

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

844 for cid in cluster_ids: 

845 self.assertLessEqual(cid, init_fake_id) 

846 # All IDs must be unique. 

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

848 

849 def testFromDagJobAttribute(self): 

850 self._copyFiles( 

851 "tiny_success", 

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

853 ) 

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

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

856 

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

858 self.assertEqual( 

859 job["from_dag_job"], 

860 "wms_tiny_success", 

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

862 ) 

863 

864 def testServiceJobPlaceholder(self): 

865 self._copyFiles( 

866 "tiny_prov_no_submit", 

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

868 ) 

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

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

871 

872 service_jobs = [ 

873 (id_, info) 

874 for id_, info in jobs.items() 

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

876 ] 

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

878 service_id, service_job = service_jobs[0] 

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

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

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

882 

883 

884class HTCJobTestCase(unittest.TestCase): 

885 """Test HTCJob methods.""" 

886 

887 def testWriteDagCommandsPayload(self): 

888 job = lssthtc.HTCJob( 

889 "job1", 

890 "label1", 

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

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

893 ) 

894 job.subfile = "job1.sub" 

895 

896 mockfh = io.StringIO() 

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

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

899 

900 def testWriteDagCommandsNotJob(self): 

901 # Testing giving command_name, no dag_rel_path and no dir 

902 job = lssthtc.HTCJob( 

903 "finalJob", 

904 "finalJob", 

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

906 ) 

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

908 mockfh = io.StringIO() 

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

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

911 

912 def testWriteDagCommandsNoop(self): 

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

914 job.subfile = "notthere.sub" 

915 mockfh = io.StringIO() 

916 job.write_dag_commands(mockfh, "") 

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

918 

919 def testWriteSubmitFile(self): 

920 job = lssthtc.HTCJob( 

921 "job1", 

922 "label1", 

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

924 ) 

925 with temporaryDirectory() as tmp_dir: 

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

927 job.write_submit_file(filename.parent) 

928 self.assertTrue(filename.exists()) 

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

930 _ = lssthtc.htc_create_submit_from_file(filename) 

931 

932 def testWriteSubmitFileExists(self): 

933 job = lssthtc.HTCJob( 

934 "job1", 

935 "label1", 

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

937 ) 

938 with temporaryDirectory() as tmp_dir: 

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

940 job.subfile = filename 

941 with open(filename, "w"): 

942 pass # make empty file 

943 job.write_submit_file(filename.parent) 

944 # make sure didn't overwrite file 

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

946 

947 

948class HtcWriteJobCommands(unittest.TestCase): 

949 """Test _htc_write_job_commands function.""" 

950 

951 def testAllCommands(self): 

952 dag_cmds = { 

953 "pre": { 

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

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

956 "executable": "exec1", 

957 "arguments": "arg1 arg2", 

958 }, 

959 "post": { 

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

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

962 "executable": "exec2", 

963 "arguments": "arg3 arg4", 

964 }, 

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

966 "pre_skip": "1", 

967 "retry": 3, 

968 "retry_unless_exit": 1, 

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

970 "priority": 123, 

971 } 

972 

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

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

975VARS job1 num="8" 

976VARS job1 spaces="a space" 

977PRE_SKIP job1 1 

978RETRY job1 3 UNLESS-EXIT 1 

979ABORT-DAG-ON job1 100 RETURN 4 

980PRIORITY job1 123 

981""" 

982 mockfh = io.StringIO() 

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

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

985 

986 def testPartialCommands(self): 

987 # Trigger skipping the inner if clauses. 

988 dag_cmds = { 

989 "pre": { 

990 "executable": "exec1", 

991 }, 

992 "post": { 

993 "executable": "exec2", 

994 }, 

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

996 "pre_skip": "1", 

997 "retry": 3, 

998 } 

999 

1000 truth = """SCRIPT PRE job1 exec1 

1001SCRIPT POST job1 exec2 

1002VARS job1 num="8" 

1003VARS job1 spaces="a space" 

1004PRE_SKIP job1 1 

1005RETRY job1 3 

1006""" 

1007 mockfh = io.StringIO() 

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

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

1010 

1011 def testNoCommands(self): 

1012 dag_cmds = {} 

1013 mockfh = io.StringIO() 

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

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

1016 

1017 def testFinal(self): 

1018 self.maxDiff = None 

1019 dag_cmds = { 

1020 "pre": { 

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

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

1023 "executable": "exec1", 

1024 "arguments": "arg1 arg2", 

1025 }, 

1026 "post": { 

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

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

1029 "executable": "exec2", 

1030 "arguments": "arg3 arg4", 

1031 }, 

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

1033 "pre_skip": "1", 

1034 "retry": 3, 

1035 "retry_unless_exit": 1, 

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

1037 "priority": 123, 

1038 } 

1039 

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

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

1042VARS finalJob num="8" 

1043VARS finalJob spaces="a space" 

1044PRE_SKIP finalJob 1 

1045""" 

1046 mockfh = io.StringIO() 

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

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

1049 

1050 

1051class HTCBackupFilesSinglePathTestCase(unittest.TestCase): 

1052 """Test htc_backup_files_single_path function.""" 

1053 

1054 def testSrcDestSame(self): 

1055 with temporaryDirectory() as tmp_dir: 

1056 with self.assertRaisesRegex( 

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

1058 ): 

1059 lssthtc.htc_backup_files_single_path(tmp_dir, tmp_dir) 

1060 

1061 def testSuccess(self): 

1062 with temporaryDirectory() as tmp_dir: 

1063 test_tmp_dir = pathlib.Path(tmp_dir) 

1064 submit_dir = test_tmp_dir / "the_src_dir" 

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

1066 backup_dir = test_tmp_dir / "the_dest_dir" 

1067 backup_dir.mkdir() 

1068 lssthtc.htc_backup_files_single_path(submit_dir, backup_dir) 

1069 result_submit = [] 

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

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

1072 self.assertEqual( 

1073 set(result_submit), 

1074 { 

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

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

1077 "./tiny_success.dag", 

1078 }, 

1079 ) 

1080 result_backup = [] 

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

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

1083 self.assertEqual( 

1084 set(result_backup), 

1085 { 

1086 "./tiny_success.info.json", 

1087 "./tiny_success.dag.metrics", 

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

1089 "./tiny_success.node_status", 

1090 }, 

1091 ) 

1092 

1093 

1094class HTCBackupFilesTestCase(unittest.TestCase): 

1095 """Test htc_backup_files function.""" 

1096 

1097 def testDirectoryNotFound(self): 

1098 with temporaryDirectory() as tmp_dir: 

1099 test_tmp_dir = pathlib.Path(tmp_dir) 

1100 submit_dir = test_tmp_dir / "submit" 

1101 with self.assertRaises(FileNotFoundError): 

1102 lssthtc.htc_backup_files(submit_dir) 

1103 

1104 def testSuccess(self): 

1105 with temporaryDirectory() as tmp_dir: 

1106 test_tmp_dir = pathlib.Path(tmp_dir) 

1107 submit_dir = test_tmp_dir / "submit" 

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

1109 lssthtc.htc_backup_files(submit_dir) 

1110 result_submit = [] 

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

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

1113 self.assertEqual( 

1114 set(result_submit), 

1115 { 

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

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

1118 "./tiny_success.dag", 

1119 "000/tiny_success.info.json", 

1120 "000/tiny_success.dag.metrics", 

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

1122 "000/tiny_success.node_status", 

1123 }, 

1124 ) 

1125 

1126 def testDestNotInSubmitDir(self): 

1127 with temporaryDirectory() as tmp_dir: 

1128 test_tmp_dir = pathlib.Path(tmp_dir) 

1129 submit_dir = test_tmp_dir / "submit" 

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

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

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

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

1134 lssthtc.htc_backup_files(submit_dir) 

1135 result_submit = [] 

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

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

1138 self.assertEqual( 

1139 set(result_submit), 

1140 { 

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

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

1143 "./tiny_problems.dag", 

1144 "./tiny_problems.dag.rescue001", 

1145 "001/tiny_problems.info.json", 

1146 "001/tiny_problems.dag.metrics", 

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

1148 "001/tiny_problems.node_status", 

1149 }, 

1150 ) 

1151 

1152 def testDestInSubmitDir(self): 

1153 with temporaryDirectory() as tmp_dir: 

1154 test_tmp_dir = pathlib.Path(tmp_dir) 

1155 submit_dir = test_tmp_dir / "submit" 

1156 backup_dir = submit_dir / "subdir" 

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

1158 lssthtc.htc_backup_files(submit_dir, backup_dir) 

1159 result_submit = [] 

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

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

1162 self.assertEqual( 

1163 set(result_submit), 

1164 { 

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

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

1167 "./tiny_problems.dag", 

1168 "./tiny_problems.dag.rescue001", 

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

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

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

1172 "subdir/001/tiny_problems.node_status", 

1173 }, 

1174 ) 

1175 

1176 def testRelativeSubdir(self): 

1177 with temporaryDirectory() as tmp_dir: 

1178 test_tmp_dir = pathlib.Path(tmp_dir) 

1179 submit_dir = test_tmp_dir / "submit" 

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

1181 lssthtc.htc_backup_files(submit_dir, "reldir") 

1182 result_submit = [] 

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

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

1185 self.assertEqual( 

1186 set(result_submit), 

1187 { 

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

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

1190 "./tiny_problems.dag", 

1191 "./tiny_problems.dag.rescue001", 

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

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

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

1195 "reldir/001/tiny_problems.node_status", 

1196 }, 

1197 ) 

1198 

1199 def testSubdags(self): 

1200 with temporaryDirectory() as tmp_dir: 

1201 test_tmp_dir = pathlib.Path(tmp_dir) 

1202 submit_dir = test_tmp_dir / "submit" 

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

1204 lssthtc.htc_backup_files(submit_dir) 

1205 result_submit = [] 

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

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

1208 self.assertEqual( 

1209 set(result_submit), 

1210 { 

1211 "./group_failed_1.dag", 

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

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

1214 "./group_failed_1.dag.rescue001", 

1215 "subdags/wms_group_order1_val1a/group_order1_val1a.dag", 

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

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

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

1219 "subdags/wms_group_order1_val1a/group_order1_val1a.node_status", 

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

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

1222 "subdags/wms_group_order1_val1b/group_order1_val1b.dag", 

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

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

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

1226 "subdags/wms_group_order1_val1c/group_order1_val1c.dag", 

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

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

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

1230 "subdags/wms_group_order1_val1c/group_order1_val1c.node_status", 

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

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

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

1234 "001/group_failed_1.info.json", 

1235 "001/group_failed_1.node_status", 

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

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

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

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

1240 }, 

1241 ) 

1242 

1243 

1244class UpdateRescueFileTestCase(unittest.TestCase): 

1245 """Test _update_rescue_file function.""" 

1246 

1247 def testSuccess(self): 

1248 self.maxDiff = None 

1249 with temporaryDirectory() as tmp_dir: 

1250 test_tmp_dir = pathlib.Path(tmp_dir) 

1251 submit_dir = test_tmp_dir / "submit" 

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

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

1254 failed_subdags = lssthtc._update_rescue_file(rescue_file) 

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

1256 with open(rescue_file) as fh: 

1257 lines = fh.readlines() 

1258 results = "".join(lines) 

1259 

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

1261# the u_testuser_DM-46294_group_fail_20250310T160455Z.dag DAG file 

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

1263# Rescue DAG version: 2.0.1 (partial) 

1264# 

1265# Total number of Nodes: 26 

1266# Nodes premarked DONE: 21 

1267# Nodes that failed: 2 

1268# wms_group_order1_val1b,finalJob,<ENDLIST> 

1269 

1270DONE pipetaskInit 

1271DONE label1_val1c_val2a 

1272DONE label1_val1b_val2b 

1273DONE label1_val1b_val2a 

1274DONE label1_val1c_val2b 

1275DONE label1_val1a_val2a 

1276DONE label1_val1a_val2b 

1277DONE label3_val1c_val2a 

1278DONE label3_val1b_val2b 

1279DONE label3_val1b_val2a 

1280DONE label3_val1c_val2b 

1281DONE label3_val1a_val2a 

1282DONE label3_val1a_val2b 

1283DONE wms_group_order1_val1a 

1284DONE label5_val1a_val2a 

1285DONE label5_val1a_val2b 

1286DONE wms_group_order1_val1c 

1287DONE label5_val1c_val2a 

1288DONE label5_val1c_val2b 

1289DONE wms_check_status_wms_group_order1_val1a 

1290DONE wms_check_status_wms_group_order1_val1c 

1291""" 

1292 

1293 self.assertEqual(results, truth) 

1294 

1295 

1296class ReadRescueHeadersTestCase(unittest.TestCase): 

1297 """Test _read_rescue_headers function.""" 

1298 

1299 def testTypical(self): 

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

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

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

1303 

1304 def testEmptyFile(self): 

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

1306 self.assertEqual(result, []) 

1307 

1308 def testOnlyHeaderLines(self): 

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

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

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

1312 

1313 def testFirstLineNotComment(self): 

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

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

1316 self.assertEqual(result, []) 

1317 

1318 def testWhitespaceStripped(self): 

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

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

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

1322 

1323 

1324class WriteRescueHeadersTestCase(unittest.TestCase): 

1325 """Test _write_rescue_headers function.""" 

1326 

1327 def testTypical(self): 

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

1329 outfh = io.StringIO() 

1330 lssthtc._write_rescue_headers(header_lines, outfh) 

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

1332 

1333 def testEmptyList(self): 

1334 outfh = io.StringIO() 

1335 lssthtc._write_rescue_headers([], outfh) 

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

1337 

1338 def testSingleLine(self): 

1339 outfh = io.StringIO() 

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

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

1342 

1343 

1344class UpdateRescueHeadersTestCase(unittest.TestCase): 

1345 """Test _update_rescue_headers function.""" 

1346 

1347 def testWithFailedSubdag(self): 

1348 header_lines = [ 

1349 "# Total number of Nodes: 26", 

1350 "# Nodes premarked DONE: 22", 

1351 "# Nodes that failed: 2", 

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

1353 ] 

1354 result = lssthtc._update_rescue_headers(header_lines) 

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

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

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

1358 

1359 def testNoSubdagFailures(self): 

1360 header_lines = [ 

1361 "# Nodes premarked DONE: 5", 

1362 "# Nodes that failed: 1", 

1363 "# finalJob,<ENDLIST>", 

1364 ] 

1365 result = lssthtc._update_rescue_headers(header_lines) 

1366 self.assertEqual(result, []) 

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

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

1369 

1370 def testMultipleFailedSubdags(self): 

1371 header_lines = [ 

1372 "# Nodes premarked DONE: 10", 

1373 "# Nodes that failed: 3", 

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

1375 ] 

1376 result = lssthtc._update_rescue_headers(header_lines) 

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

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

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

1380 

1381 def testNoFailedNodesLine(self): 

1382 header_lines = [ 

1383 "# Total number of Nodes: 5", 

1384 "# Nodes premarked DONE: 5", 

1385 ] 

1386 original = list(header_lines) 

1387 result = lssthtc._update_rescue_headers(header_lines) 

1388 self.assertEqual(result, []) 

1389 self.assertEqual(header_lines, original) 

1390 

1391 def testEmptyHeader(self): 

1392 result = lssthtc._update_rescue_headers([]) 

1393 self.assertEqual(result, []) 

1394 

1395 

1396class ReadDagStatusTestCase(unittest.TestCase): 

1397 """Test read_dag_status function and read_single_dag_status.""" 

1398 

1399 def testFileMissing(self): 

1400 with temporaryDirectory() as tmp_dir: 

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

1402 _ = lssthtc.read_dag_status(tmp_dir) 

1403 

1404 def testRegular(self): 

1405 with temporaryDirectory() as tmp_dir: 

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

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

1408 results = lssthtc.read_dag_status(submit_dir) 

1409 truth = { 

1410 "JobProcsHeld": 0, 

1411 "NodesPost": 0, 

1412 "JobProcsIdle": 0, 

1413 "NodesTotal": 6, 

1414 "NodesFailed": 2, 

1415 "NodesDone": 3, 

1416 "NodesQueued": 0, 

1417 "NodesPre": 0, 

1418 "NodesFutile": 1, 

1419 "NodesUnready": 0, 

1420 } 

1421 self.assertEqual(results, results | truth) 

1422 

1423 def testSubdags(self): 

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

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

1426 """ 

1427 self.maxDiff = None 

1428 with temporaryDirectory() as tmp_dir: 

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

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

1431 results = lssthtc.read_dag_status(submit_dir) 

1432 truth = { 

1433 "JobProcsHeld": 0, 

1434 "NodesPost": 0, 

1435 "JobProcsIdle": 0, 

1436 "NodesTotal": 34, 

1437 "NodesFailed": 0, 

1438 "NodesDone": 17, 

1439 "NodesQueued": 3, 

1440 "NodesPre": 0, 

1441 "NodesFutile": 0, 

1442 "NodesUnready": 14, 

1443 } 

1444 self.assertEqual(results, results | truth) 

1445 

1446 

1447class ReadDagInfoTestCase(unittest.TestCase): 

1448 """Test read_dag_info function.""" 

1449 

1450 def testFileMissing(self): 

1451 with temporaryDirectory() as tmp_dir: 

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

1453 _ = lssthtc.read_dag_info(tmp_dir) 

1454 

1455 def testSuccess(self): 

1456 with temporaryDirectory() as tmp_dir: 

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

1458 filename, results = lssthtc.read_dag_info(tmp_dir) 

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

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

1461 

1462 truth = { 

1463 "test02": { 

1464 "9208.0": { 

1465 "ClusterId": 9208, 

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

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

1468 "bps_project": "dev", 

1469 "bps_payload": "tiny", 

1470 "bps_operator": "testuser", 

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

1472 "bps_provisioning_job": "provisioningJob", 

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

1474 "bps_campaign": "quick", 

1475 "bps_runsite": "testpool", 

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

1477 "bps_run": "u_testuser_tiny_20250213T164427Z", 

1478 "bps_isjob": "True", 

1479 } 

1480 } 

1481 } 

1482 

1483 self.assertEqual(results, truth) 

1484 

1485 def testPermissionError(self): 

1486 with temporaryDirectory() as tmp_dir: 

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

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

1489 mocked_open.side_effect = PermissionError 

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

1491 _, results = lssthtc.read_dag_info(tmp_dir) 

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

1493 self.assertEqual({}, results) 

1494 

1495 

1496class HtcWriteCondorFileTestCase(unittest.TestCase): 

1497 """Test htc_write_condor_file function.""" 

1498 

1499 def testSuccess(self): 

1500 with temporaryDirectory() as tmp_dir: 

1501 job_name = "job1" 

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

1503 job = { 

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

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

1506 "request_memory": "2000", 

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

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

1509 } 

1510 job_attrs = { 

1511 "bps_job_name": job_name, 

1512 "bps_job_label": "label1", 

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

1514 } 

1515 expected = [ 

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

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

1518 "request_memory=2000\n", 

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

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

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

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

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

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

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

1526 "queue\n", 

1527 ] 

1528 

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

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

1531 actual = f.readlines() 

1532 

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

1534 self.assertTrue(filename.exists()) 

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

1536 _ = lssthtc.htc_create_submit_from_file(filename) 

1537 

1538 

1539class HtcCreateSubmitFromDagTestCase(unittest.TestCase): 

1540 """Test htc_create_submit_from_dag function.""" 

1541 

1542 @classmethod 

1543 def setUpClass(cls): 

1544 cls.bindir = None 

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

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

1547 cls.bindir = tempfile.TemporaryDirectory() 

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

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

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

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

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

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

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

1555 

1556 @classmethod 

1557 def tearDownClass(cls): 

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

1559 cls.bindir.cleanup() 

1560 

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

1562 def testMaxIdleEnvVar(self): 

1563 with temporaryDirectory() as tmp_dir: 

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

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

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

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

1568 

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

1570 def testMaxIdleInDAGManConfig(self): 

1571 with temporaryDirectory() as tmp_dir: 

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

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

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

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

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

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

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

1579 

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

1581 def testMaxIdleNotInDAGManConfig(self): 

1582 with temporaryDirectory() as tmp_dir: 

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

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

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

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

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

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

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

1590 

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

1592 def testMaxIdleGiven(self): 

1593 with temporaryDirectory() as tmp_dir: 

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

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

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

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

1598 

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

1600 def testMaxJobsIdleParam(self): 

1601 def _fake_params_contains(key): 

1602 if key == "DAGMAN_MAX_JOBS_IDLE": 

1603 return True 

1604 return False # pragma: no cover 

1605 

1606 def _fake_params_get(key): 

1607 if key == "DAGMAN_MAX_JOBS_IDLE": 

1608 return 16 

1609 return "FAKE_VAL" # pragma: no cover 

1610 

1611 with temporaryDirectory() as tmp_dir: 

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

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

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

1615 mock_param.__contains__.side_effect = _fake_params_contains 

1616 mock_param.__getitem__.side_effect = _fake_params_get 

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

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

1619 

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

1621 def testNoMaxJobsIdle(self): 

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

1623 HTCondor version when no MaxIdle passed to from_dag, not 

1624 checking arguments string here. Instead just making sure 

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

1626 """ 

1627 with temporaryDirectory() as tmp_dir: 

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

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

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

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

1632 mock_param.__contains__.return_value = False 

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

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

1635 

1636 

1637class HtcDagTestCase(unittest.TestCase): 

1638 """Test for HTCDag class.""" 

1639 

1640 def setUp(self): 

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

1642 job.add_job_cmds( 

1643 { 

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

1645 "arguments": "foo", 

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

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

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

1649 } 

1650 ) 

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

1652 

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

1654 self.dag.add_job(job) 

1655 

1656 self.subfile_expected = [ 

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

1658 'arguments="foo"\n', 

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

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

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

1662 "queue\n", 

1663 ] 

1664 

1665 def tearDown(self): 

1666 pass 

1667 

1668 def testWriteWithDagConfig(self): 

1669 with temporaryDirectory() as tmp_dir: 

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

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

1672 wms_config_filename = "dagman.conf" 

1673 wms_configurator = dagman_configurator.DagmanConfigurator(config) 

1674 wms_configurator.prepare(wms_config_filename, prefix=tmp_dir) 

1675 wms_configurator.configure(self.dag) 

1676 dagfile_expected = [ 

1677 f"CONFIG {wms_config_filename}\n", 

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

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

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

1681 ] 

1682 

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

1684 

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

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

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

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

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

1690 dagfile_actual = f.readlines() 

1691 self.assertEqual(dagfile_actual, dagfile_expected) 

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

1693 subfile_actual = f.readlines() 

1694 self.assertEqual(subfile_actual, self.subfile_expected) 

1695 

1696 def testWriteWithoutDagConfig(self): 

1697 with temporaryDirectory() as tmp_dir: 

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

1699 dagfile_expected = [ 

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

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

1702 ] 

1703 

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

1705 

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

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

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

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

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

1711 dagfile_actual = f.readlines() 

1712 self.assertEqual(dagfile_actual, dagfile_expected) 

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

1714 subfile_actual = f.readlines() 

1715 self.assertEqual(subfile_actual, self.subfile_expected) 

1716 

1717 def testWriteLazySubdag(self): 

1718 self.maxDiff = None 

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

1720 dag.graph["write_dot"] = True 

1721 with temporaryDirectory() as tmp_dir: 

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

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

1724 dagfile_actual = f.readlines() 

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

1726 

1727 all_files = [] 

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

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

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

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

1732 

1733 @staticmethod 

1734 def _make_simple_job(name): 

1735 job = lssthtc.HTCJob(name=name) 

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

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

1738 return job 

1739 

1740 def testWriteEdgesAndSpecialJobs(self): 

1741 self.maxDiff = None 

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

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

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

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

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

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

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

1749 

1750 dagfile_expected = [ 

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

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

1753 "PARENT jobA CHILD jobB\n", 

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

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

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

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

1758 ] 

1759 

1760 with temporaryDirectory() as tmp_dir: 

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

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

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

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

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

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

1767 

1768 def testWriteDagSubdir(self): 

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

1770 # that subdirectory, creating it as needed. 

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

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

1773 

1774 with temporaryDirectory() as tmp_dir: 

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

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

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

1778 

1779 def testWriteMissingJobData(self): 

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

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

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

1783 with temporaryDirectory() as tmp_dir: 

1784 with self.assertRaises(KeyError): 

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

1786 

1787 

1788class WriteDagInfoTestCase(unittest.TestCase): 

1789 """Test for write_dag_info function.""" 

1790 

1791 def setUp(self): 

1792 self.run = "u_testuser_DM-53494_20260220T001651Z" 

1793 self.data = { 

1794 "mycomputer": { 

1795 "24390.0": { 

1796 "ClusterId": 24390, 

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

1798 "bps_run": self.run, 

1799 "bps_isjob": "True", 

1800 "bps_payload": "DM-53494", 

1801 "bps_project": "dev", 

1802 "bps_runsite": "site1", 

1803 "bps_campaign": "ci_rc2", 

1804 "bps_operator": "testuser", 

1805 "bps_run_quanta": "", 

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

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

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

1809 "bps_wms_config_path": "dagman.conf", 

1810 } 

1811 } 

1812 } 

1813 

1814 def testWrite(self): 

1815 with temporaryDirectory() as tmp_dir: 

1816 with chdir(tmp_dir): 

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

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

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

1820 self.assertEqual(filename, path) 

1821 

1822 read_filename, read_data = lssthtc.read_dag_info(tmp_dir) 

1823 self.assertEqual(read_filename, path) 

1824 self.assertEqual(read_data, self.data) 

1825 

1826 

1827if __name__ == "__main__": 

1828 unittest.main()