Coverage for tests/test_lssthtc.py: 99%

761 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-07-22 17:25 +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 

39 

40from lsst.ctrl.bps import BpsConfig 

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

42from lsst.daf.butler import Config 

43from lsst.utils.tests import temporaryDirectory 

44 

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

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

47 

48 

49class TestLsstHtc(unittest.TestCase): 

50 """Test basic usage.""" 

51 

52 def testHtcEscapeInt(self): 

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

54 

55 def testHtcEscapeDouble(self): 

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

57 

58 def testHtcEscapeSingle(self): 

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

60 

61 def testHtcEscapeNoSideEffect(self): 

62 val = "'val'" 

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

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

65 

66 def testHtcEscapeQuot(self): 

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

68 

69 def testHtcVersion(self): 

70 ver = lssthtc.htc_version() 

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

72 

73 

74class HtcTweakJobInfoTestCase(unittest.TestCase): 

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

76 

77 def setUp(self): 

78 self.log_dir = tempfile.TemporaryDirectory() 

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

80 self.job = { 

81 "Cluster": 1, 

82 "Proc": 0, 

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

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

85 "MyType": None, 

86 "TerminatedNormally": True, 

87 } 

88 

89 def tearDown(self): 

90 self.log_dir.cleanup() 

91 

92 def testDirectAssignments(self): 

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

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

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

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

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

98 

99 def testIncompatibleAdPassThru(self): 

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

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

102 result = dict(expected) 

103 lssthtc.htc_tweak_log_info(self.log_dirname, result) 

104 self.assertEqual(result, expected) 

105 

106 def testJobStatusAssignmentJobAbortedEvent(self): 

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

108 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

109 self.assertTrue("JobStatus" in job) 

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

111 

112 def testJobStatusAssignmentExecuteEvent(self): 

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

114 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

115 self.assertTrue("JobStatus" in job) 

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

117 

118 def testJobStatusAssignmentSubmitEvent(self): 

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

120 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

121 self.assertTrue("JobStatus" in job) 

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

123 

124 def testJobStatusAssignmentJobHeldEvent(self): 

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

126 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

127 self.assertTrue("JobStatus" in job) 

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

129 

130 def testJobStatusAssignmentJobTerminatedEvent(self): 

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

132 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

133 self.assertTrue("JobStatus" in job) 

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

135 

136 def testJobStatusAssignmentPostScriptTerminatedEvent(self): 

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

138 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

139 self.assertTrue("JobStatus" in job) 

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

141 

142 def testJobStatusAssignmentReleaseEventMainDagJob(self): 

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

144 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

145 self.assertTrue("JobStatus" in job) 

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

147 

148 def testJobStatusAssignmentReleaseEventForNodeJob(self): 

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

150 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

151 self.assertTrue("JobStatus" in job) 

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

153 

154 def testAddingExitStatusSuccess(self): 

155 job = self.job | { 

156 "MyType": "JobTerminatedEvent", 

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

158 } 

159 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

160 self.assertIn("ExitBySignal", job) 

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

162 self.assertIn("ExitCode", job) 

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

164 

165 def testAddingExitStatusFailure(self): 

166 job = self.job | { 

167 "MyType": "JobHeldEvent", 

168 } 

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

170 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

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

172 

173 def testLoggingUnknownLogEvent(self): 

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

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

176 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

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

178 

179 def testMissingKey(self): 

180 job = self.job 

181 del job["Cluster"] 

182 with self.assertRaises(KeyError) as cm: 

183 lssthtc.htc_tweak_log_info(self.log_dirname, job) 

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

185 

186 

187class HtcCheckDagmanOutputTestCase(unittest.TestCase): 

188 """Test htc_check_dagman_output function.""" 

189 

190 def test_missing_output_file(self): 

191 with temporaryDirectory() as tmp_dir: 

192 with self.assertRaises(FileNotFoundError): 

193 _ = lssthtc.htc_check_dagman_output(tmp_dir) 

194 

195 def test_permissions_output_file(self): 

196 with temporaryDirectory() as tmp_dir: 

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

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

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

200 results = lssthtc.htc_check_dagman_output(tmp_dir) 

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

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

203 

204 def test_submit_failure(self): 

205 with temporaryDirectory() as tmp_dir: 

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

207 results = lssthtc.htc_check_dagman_output(tmp_dir) 

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

209 

210 def test_tmpdir_abort(self): 

211 with temporaryDirectory() as tmp_dir: 

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

213 results = lssthtc.htc_check_dagman_output(tmp_dir) 

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

215 

216 def test_no_messages(self): 

217 with temporaryDirectory() as tmp_dir: 

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

219 results = lssthtc.htc_check_dagman_output(tmp_dir) 

220 self.assertEqual("", results) 

221 

222 

223class SummarizeDagTestCase(unittest.TestCase): 

224 """Test summarize_dag function.""" 

225 

226 def test_no_dag_file(self): 

227 with temporaryDirectory() as tmp_dir: 

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

229 self.assertFalse(len(job_name_to_pipetask)) 

230 self.assertFalse(len(job_name_to_type)) 

231 self.assertFalse(summary) 

232 

233 def test_success(self): 

234 with temporaryDirectory() as tmp_dir: 

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

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

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

238 self.assertEqual( 

239 job_name_to_label, 

240 { 

241 "pipetaskInit": "pipetaskInit", 

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

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

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

245 "finalJob": "finalJob", 

246 }, 

247 ) 

248 self.assertEqual( 

249 job_name_to_type, 

250 { 

251 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

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

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

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

255 "finalJob": lssthtc.WmsNodeType.FINAL, 

256 }, 

257 ) 

258 

259 def test_service(self): 

260 with temporaryDirectory() as tmp_dir: 

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

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

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

264 self.assertEqual( 

265 job_name_to_label, 

266 { 

267 "pipetaskInit": "pipetaskInit", 

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

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

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

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

272 "finalJob": "finalJob", 

273 "provisioningJob": "provisioningJob", 

274 }, 

275 ) 

276 self.assertEqual( 

277 job_name_to_type, 

278 { 

279 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

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

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

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

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

284 "finalJob": lssthtc.WmsNodeType.FINAL, 

285 "provisioningJob": lssthtc.WmsNodeType.SERVICE, 

286 }, 

287 ) 

288 

289 def test_noop(self): 

290 with temporaryDirectory() as tmp_dir: 

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

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

293 self.assertEqual( 

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

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

296 ) 

297 self.assertEqual( 

298 job_name_to_label, 

299 { 

300 "label1_val1a_val2a": "label1", 

301 "label1_val1a_val2b": "label1", 

302 "label1_val1b_val2a": "label1", 

303 "label1_val1b_val2b": "label1", 

304 "label1_val1c_val2a": "label1", 

305 "label1_val1c_val2b": "label1", 

306 "label2_val1a_val2a": "label2", 

307 "label2_val1a_val2b": "label2", 

308 "label2_val1b_val2a": "label2", 

309 "label2_val1b_val2b": "label2", 

310 "label2_val1c_val2a": "label2", 

311 "label2_val1c_val2b": "label2", 

312 "label3_val1a_val2a": "label3", 

313 "label3_val1a_val2b": "label3", 

314 "label3_val1b_val2a": "label3", 

315 "label3_val1b_val2b": "label3", 

316 "label3_val1c_val2a": "label3", 

317 "label3_val1c_val2b": "label3", 

318 "label4_val1a_val2a": "label4", 

319 "label4_val1a_val2b": "label4", 

320 "label4_val1b_val2a": "label4", 

321 "label4_val1b_val2b": "label4", 

322 "label4_val1c_val2a": "label4", 

323 "label4_val1c_val2b": "label4", 

324 "label5_val1a_val2a": "label5", 

325 "label5_val1a_val2b": "label5", 

326 "label5_val1b_val2a": "label5", 

327 "label5_val1b_val2b": "label5", 

328 "label5_val1c_val2a": "label5", 

329 "label5_val1c_val2b": "label5", 

330 "finalJob": "finalJob", 

331 "pipetaskInit": "pipetaskInit", 

332 "wms_noop_order1_val1a": "order1", 

333 "wms_noop_order1_val1b": "order1", 

334 }, 

335 ) 

336 self.assertEqual( 

337 job_name_to_type, 

338 { 

339 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

340 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

341 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

342 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

343 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

344 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

345 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

346 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

347 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

348 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

349 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

350 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

351 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

352 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

353 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

354 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

355 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

356 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

357 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

358 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

359 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

360 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

361 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

362 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

363 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

364 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

365 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

366 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

367 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

368 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

369 "finalJob": lssthtc.WmsNodeType.FINAL, 

370 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

371 "wms_noop_order1_val1a": lssthtc.WmsNodeType.NOOP, 

372 "wms_noop_order1_val1b": lssthtc.WmsNodeType.NOOP, 

373 }, 

374 ) 

375 

376 def test_subdags(self): 

377 with temporaryDirectory() as tmp_dir: 

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

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

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

381 self.assertEqual( 

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

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

384 ) 

385 

386 self.assertEqual( 

387 job_name_to_label, 

388 { 

389 "pipetaskInit": "pipetaskInit", 

390 "label1_val1b_val2a": "label1", 

391 "label1_val1c_val2a": "label1", 

392 "label1_val1a_val2b": "label1", 

393 "label1_val1b_val2b": "label1", 

394 "label1_val1c_val2b": "label1", 

395 "label1_val1a_val2a": "label1", 

396 "label2_val1a_val2b": "label2", 

397 "label2_val1a_val2a": "label2", 

398 "label2_val1b_val2a": "label2", 

399 "label2_val1b_val2b": "label2", 

400 "label2_val1c_val2a": "label2", 

401 "label2_val1c_val2b": "label2", 

402 "label3_val1b_val2a": "label3", 

403 "label3_val1c_val2a": "label3", 

404 "label3_val1a_val2b": "label3", 

405 "label3_val1b_val2b": "label3", 

406 "label3_val1c_val2b": "label3", 

407 "label3_val1a_val2a": "label3", 

408 "label4_val1a_val2b": "label4", 

409 "label4_val1a_val2a": "label4", 

410 "label4_val1b_val2a": "label4", 

411 "label4_val1b_val2b": "label4", 

412 "label4_val1c_val2a": "label4", 

413 "label4_val1c_val2b": "label4", 

414 "label5_val1a_val2b": "label5", 

415 "label5_val1a_val2a": "label5", 

416 "label5_val1b_val2a": "label5", 

417 "label5_val1b_val2b": "label5", 

418 "label5_val1c_val2a": "label5", 

419 "label5_val1c_val2b": "label5", 

420 "finalJob": "finalJob", 

421 "provisioningJob": "provisioningJob", 

422 "wms_group_order1_val1a": "order1", 

423 "wms_group_order1_val1b": "order1", 

424 "wms_group_order1_val1c": "order1", 

425 "wms_check_status_wms_group_order1_val1a": "order1", 

426 "wms_check_status_wms_group_order1_val1b": "order1", 

427 "wms_check_status_wms_group_order1_val1c": "order1", 

428 }, 

429 ) 

430 

431 self.assertEqual( 

432 job_name_to_type, 

433 { 

434 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD, 

435 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

436 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

437 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

438 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

439 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

440 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

441 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

442 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

443 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

444 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

445 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

446 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

447 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

448 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

449 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

450 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

451 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

452 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

453 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

454 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

455 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

456 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

457 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

458 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

459 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD, 

460 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD, 

461 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD, 

462 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD, 

463 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD, 

464 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD, 

465 "finalJob": lssthtc.WmsNodeType.FINAL, 

466 "provisioningJob": lssthtc.WmsNodeType.SERVICE, 

467 "wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG, 

468 "wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG, 

469 "wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG, 

470 "wms_check_status_wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG_CHECK, 

471 "wms_check_status_wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG_CHECK, 

472 "wms_check_status_wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG_CHECK, 

473 }, 

474 ) 

475 

476 

477class ReadDagNodesLogTestCase(unittest.TestCase): 

478 """Test read_dag_nodes_log function.""" 

479 

480 def setUp(self): 

481 self.tmpdir = tempfile.mkdtemp() 

482 

483 def tearDown(self): 

484 rmtree(self.tmpdir, ignore_errors=True) 

485 

486 def testFileMissing(self): 

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

488 _ = lssthtc.read_dag_nodes_log(self.tmpdir) 

489 

490 def testRegular(self): 

491 with temporaryDirectory() as tmp_dir: 

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

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

494 results = lssthtc.read_dag_nodes_log(submit_dir) 

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

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

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

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

499 

500 def testSubdags(self): 

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

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

503 """ 

504 with temporaryDirectory() as tmp_dir: 

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

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

507 results = lssthtc.read_dag_nodes_log(submit_dir) 

508 # main dag 

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

510 # subdag 

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

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

513 

514 

515class ReadNodeStatusTestCase(unittest.TestCase): 

516 """Test read_node_status function.""" 

517 

518 def setUp(self): 

519 self.tmpdir = tempfile.mkdtemp() 

520 

521 def tearDown(self): 

522 rmtree(self.tmpdir, ignore_errors=True) 

523 

524 def testServiceJobNotSubmitted(self): 

525 # tiny_prov_no_submit files have successful workflow 

526 # but provisioningJob could not submit. 

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

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

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

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

531 

532 jobs = lssthtc.read_node_status(self.tmpdir) 

533 found = [ 

534 id_ 

535 for id_ in jobs 

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

537 ] 

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

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

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

541 

542 def testMissingStatusFile(self): 

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

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

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

546 

547 jobs = lssthtc.read_node_status(self.tmpdir) 

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

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

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

551 found = [ 

552 id_ 

553 for id_ in jobs 

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

555 ] 

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

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

558 

559 def testSubdagsRunning(self): 

560 with temporaryDirectory() as tmp_dir: 

561 test_tmp_dir = pathlib.Path(tmp_dir) 

562 submit_dir = test_tmp_dir / "submit" 

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

564 jobs = lssthtc.read_node_status(submit_dir) 

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

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

567 job_name_to_id = {} 

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

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

570 job_type_to_names = {} 

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

572 job_type_to_names.setdefault( 

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

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

575 

576 # check counts 

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

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

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

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

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

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

583 

584 # spot check some statuses 

585 self.assertEqual( 

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

587 ) 

588 self.assertEqual( 

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

590 ) 

591 self.assertEqual( 

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

593 ) 

594 self.assertEqual( 

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

596 ) 

597 

598 def testSubdagsFailed(self): 

599 with temporaryDirectory() as tmp_dir: 

600 test_tmp_dir = pathlib.Path(tmp_dir) 

601 submit_dir = test_tmp_dir / "submit" 

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

603 jobs = lssthtc.read_node_status(submit_dir) 

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

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

606 job_name_to_id = {} 

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

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

609 job_type_to_names = {} 

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

611 job_type_to_names.setdefault( 

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

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

614 

615 # check counts 

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

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

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

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

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

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

622 

623 # spot check some statuses 

624 self.assertEqual( 

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

626 ) 

627 self.assertEqual( 

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

629 ) 

630 self.assertEqual( 

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

632 ) 

633 

634 self.assertEqual( 

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

636 ) 

637 self.assertEqual( 

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

639 ) 

640 self.assertEqual( 

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

642 lssthtc.NodeStatus.ERROR, 

643 ) 

644 

645 

646class ReadSingleNodeStatusTestCase(unittest.TestCase): 

647 """Test read_single_node_status function.""" 

648 

649 def setUp(self): 

650 self.tmpdir = tempfile.mkdtemp() 

651 

652 def tearDown(self): 

653 rmtree(self.tmpdir, ignore_errors=True) 

654 

655 def _copyFiles(self, data_subdir, suffixes): 

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

657 for suffix in suffixes: 

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

659 

660 def _jobNameToId(self, jobs): 

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

662 

663 def testAllDone(self): 

664 self._copyFiles( 

665 "tiny_success", 

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

667 ) 

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

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

670 

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

672 name_to_id = self._jobNameToId(jobs) 

673 

674 # All four submitted nodes are marked DONE. 

675 for name in [ 

676 "pipetaskInit", 

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

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

679 "finalJob", 

680 ]: 

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

682 self.assertEqual( 

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

684 lssthtc.NodeStatus.DONE, 

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

686 ) 

687 

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

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

690 self.assertIn("provisioningJob", name_to_id) 

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

692 

693 # Spot-check labels and types. 

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

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

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

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

698 

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

700 for job in jobs.values(): 

701 self.assertIn("DAGManJobID", job) 

702 

703 def testMixedStatuses(self): 

704 self._copyFiles( 

705 "tiny_problems", 

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

707 ) 

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

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

710 

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

712 name_to_id = self._jobNameToId(jobs) 

713 

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

715 self.assertEqual( 

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

717 lssthtc.NodeStatus.ERROR, 

718 ) 

719 self.assertEqual( 

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

721 lssthtc.NodeStatus.DONE, 

722 ) 

723 self.assertEqual( 

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

725 lssthtc.NodeStatus.FUTILE, 

726 ) 

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

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

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

730 self.assertIn("provisioningJob", name_to_id) 

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

732 

733 def testRunningWorkflow(self): 

734 self._copyFiles( 

735 "tiny_running", 

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

737 ) 

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

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

740 

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

742 name_to_id = self._jobNameToId(jobs) 

743 

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

745 self.assertEqual( 

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

747 lssthtc.NodeStatus.SUBMITTED, 

748 ) 

749 self.assertEqual( 

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

751 lssthtc.NodeStatus.NOT_READY, 

752 ) 

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

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

755 self.assertIn("provisioningJob", name_to_id) 

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

757 

758 def testMissingNodeStatusFile(self): 

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

760 # and dag. 

761 self._copyFiles( 

762 "tiny_problems", 

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

764 ) 

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

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

767 

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

769 name_to_id = self._jobNameToId(jobs) 

770 

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

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

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

774 

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

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

777 

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

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

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

781 

782 def testMissingLogFiles(self): 

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

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

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

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

787 

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

789 for job in jobs.values(): 

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

791 

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

793 name_to_id = self._jobNameToId(jobs) 

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

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

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

797 

798 def testInitFakeId(self): 

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

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

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

802 init_fake_id = -10 

803 jobs = lssthtc.read_single_node_status(filename, init_fake_id) 

804 

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

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

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

808 for cid in cluster_ids: 

809 self.assertLessEqual(cid, init_fake_id) 

810 # All IDs must be unique. 

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

812 

813 def testFromDagJobAttribute(self): 

814 self._copyFiles( 

815 "tiny_success", 

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

817 ) 

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

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

820 

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

822 self.assertEqual( 

823 job["from_dag_job"], 

824 "wms_tiny_success", 

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

826 ) 

827 

828 def testServiceJobPlaceholder(self): 

829 self._copyFiles( 

830 "tiny_prov_no_submit", 

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

832 ) 

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

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

835 

836 service_jobs = [ 

837 (id_, info) 

838 for id_, info in jobs.items() 

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

840 ] 

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

842 service_id, service_job = service_jobs[0] 

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

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

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

846 

847 

848class HTCJobTestCase(unittest.TestCase): 

849 """Test HTCJob methods.""" 

850 

851 def testWriteDagCommandsPayload(self): 

852 job = lssthtc.HTCJob( 

853 "job1", 

854 "label1", 

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

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

857 ) 

858 job.subfile = "job1.sub" 

859 

860 mockfh = io.StringIO() 

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

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

863 

864 def testWriteDagCommandsNotJob(self): 

865 # Testing giving command_name, no dag_rel_path and no dir 

866 job = lssthtc.HTCJob( 

867 "finalJob", 

868 "finalJob", 

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

870 ) 

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

872 mockfh = io.StringIO() 

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

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

875 

876 def testWriteDagCommandsNoop(self): 

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

878 job.subfile = "notthere.sub" 

879 mockfh = io.StringIO() 

880 job.write_dag_commands(mockfh, "") 

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

882 

883 def testWriteSubmitFile(self): 

884 job = lssthtc.HTCJob( 

885 "job1", 

886 "label1", 

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

888 ) 

889 with temporaryDirectory() as tmp_dir: 

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

891 job.write_submit_file(filename.parent) 

892 self.assertTrue(filename.exists()) 

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

894 _ = lssthtc.htc_create_submit_from_file(filename) 

895 

896 def testWriteSubmitFileExists(self): 

897 job = lssthtc.HTCJob( 

898 "job1", 

899 "label1", 

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

901 ) 

902 with temporaryDirectory() as tmp_dir: 

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

904 job.subfile = filename 

905 with open(filename, "w"): 

906 pass # make empty file 

907 job.write_submit_file(filename.parent) 

908 # make sure didn't overwrite file 

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

910 

911 

912class HtcWriteJobCommands(unittest.TestCase): 

913 """Test _htc_write_job_commands function.""" 

914 

915 def testAllCommands(self): 

916 dag_cmds = { 

917 "pre": { 

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

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

920 "executable": "exec1", 

921 "arguments": "arg1 arg2", 

922 }, 

923 "post": { 

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

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

926 "executable": "exec2", 

927 "arguments": "arg3 arg4", 

928 }, 

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

930 "pre_skip": "1", 

931 "retry": 3, 

932 "retry_unless_exit": 1, 

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

934 "priority": 123, 

935 } 

936 

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

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

939VARS job1 num="8" 

940VARS job1 spaces="a space" 

941PRE_SKIP job1 1 

942RETRY job1 3 UNLESS-EXIT 1 

943ABORT-DAG-ON job1 100 RETURN 4 

944PRIORITY job1 123 

945""" 

946 mockfh = io.StringIO() 

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

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

949 

950 def testPartialCommands(self): 

951 # Trigger skipping the inner if clauses. 

952 dag_cmds = { 

953 "pre": { 

954 "executable": "exec1", 

955 }, 

956 "post": { 

957 "executable": "exec2", 

958 }, 

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

960 "pre_skip": "1", 

961 "retry": 3, 

962 } 

963 

964 truth = """SCRIPT PRE job1 exec1 

965SCRIPT POST job1 exec2 

966VARS job1 num="8" 

967VARS job1 spaces="a space" 

968PRE_SKIP job1 1 

969RETRY job1 3 

970""" 

971 mockfh = io.StringIO() 

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

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

974 

975 def testNoCommands(self): 

976 dag_cmds = {} 

977 mockfh = io.StringIO() 

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

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

980 

981 def testFinal(self): 

982 self.maxDiff = None 

983 dag_cmds = { 

984 "pre": { 

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

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

987 "executable": "exec1", 

988 "arguments": "arg1 arg2", 

989 }, 

990 "post": { 

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

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

993 "executable": "exec2", 

994 "arguments": "arg3 arg4", 

995 }, 

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

997 "pre_skip": "1", 

998 "retry": 3, 

999 "retry_unless_exit": 1, 

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

1001 "priority": 123, 

1002 } 

1003 

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

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

1006VARS finalJob num="8" 

1007VARS finalJob spaces="a space" 

1008PRE_SKIP finalJob 1 

1009""" 

1010 mockfh = io.StringIO() 

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

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

1013 

1014 

1015class HTCBackupFilesSinglePathTestCase(unittest.TestCase): 

1016 """Test htc_backup_files_single_path function.""" 

1017 

1018 def testSrcDestSame(self): 

1019 with temporaryDirectory() as tmp_dir: 

1020 with self.assertRaisesRegex( 

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

1022 ): 

1023 lssthtc.htc_backup_files_single_path(tmp_dir, tmp_dir) 

1024 

1025 def testSuccess(self): 

1026 with temporaryDirectory() as tmp_dir: 

1027 test_tmp_dir = pathlib.Path(tmp_dir) 

1028 submit_dir = test_tmp_dir / "the_src_dir" 

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

1030 backup_dir = test_tmp_dir / "the_dest_dir" 

1031 backup_dir.mkdir() 

1032 lssthtc.htc_backup_files_single_path(submit_dir, backup_dir) 

1033 result_submit = [] 

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

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

1036 self.assertEqual( 

1037 set(result_submit), 

1038 { 

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

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

1041 "./tiny_success.dag", 

1042 }, 

1043 ) 

1044 result_backup = [] 

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

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

1047 self.assertEqual( 

1048 set(result_backup), 

1049 { 

1050 "./tiny_success.info.json", 

1051 "./tiny_success.dag.metrics", 

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

1053 "./tiny_success.node_status", 

1054 }, 

1055 ) 

1056 

1057 

1058class HTCBackupFilesTestCase(unittest.TestCase): 

1059 """Test htc_backup_files function.""" 

1060 

1061 def testDirectoryNotFound(self): 

1062 with temporaryDirectory() as tmp_dir: 

1063 test_tmp_dir = pathlib.Path(tmp_dir) 

1064 submit_dir = test_tmp_dir / "submit" 

1065 with self.assertRaises(FileNotFoundError): 

1066 lssthtc.htc_backup_files(submit_dir) 

1067 

1068 def testSuccess(self): 

1069 with temporaryDirectory() as tmp_dir: 

1070 test_tmp_dir = pathlib.Path(tmp_dir) 

1071 submit_dir = test_tmp_dir / "submit" 

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

1073 lssthtc.htc_backup_files(submit_dir) 

1074 result_submit = [] 

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

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

1077 self.assertEqual( 

1078 set(result_submit), 

1079 { 

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

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

1082 "./tiny_success.dag", 

1083 "000/tiny_success.info.json", 

1084 "000/tiny_success.dag.metrics", 

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

1086 "000/tiny_success.node_status", 

1087 }, 

1088 ) 

1089 

1090 def testDestNotInSubmitDir(self): 

1091 with temporaryDirectory() as tmp_dir: 

1092 test_tmp_dir = pathlib.Path(tmp_dir) 

1093 submit_dir = test_tmp_dir / "submit" 

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

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

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

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

1098 lssthtc.htc_backup_files(submit_dir) 

1099 result_submit = [] 

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

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

1102 self.assertEqual( 

1103 set(result_submit), 

1104 { 

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

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

1107 "./tiny_problems.dag", 

1108 "./tiny_problems.dag.rescue001", 

1109 "001/tiny_problems.info.json", 

1110 "001/tiny_problems.dag.metrics", 

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

1112 "001/tiny_problems.node_status", 

1113 }, 

1114 ) 

1115 

1116 def testDestInSubmitDir(self): 

1117 with temporaryDirectory() as tmp_dir: 

1118 test_tmp_dir = pathlib.Path(tmp_dir) 

1119 submit_dir = test_tmp_dir / "submit" 

1120 backup_dir = submit_dir / "subdir" 

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

1122 lssthtc.htc_backup_files(submit_dir, backup_dir) 

1123 result_submit = [] 

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

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

1126 self.assertEqual( 

1127 set(result_submit), 

1128 { 

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

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

1131 "./tiny_problems.dag", 

1132 "./tiny_problems.dag.rescue001", 

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

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

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

1136 "subdir/001/tiny_problems.node_status", 

1137 }, 

1138 ) 

1139 

1140 def testRelativeSubdir(self): 

1141 with temporaryDirectory() as tmp_dir: 

1142 test_tmp_dir = pathlib.Path(tmp_dir) 

1143 submit_dir = test_tmp_dir / "submit" 

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

1145 lssthtc.htc_backup_files(submit_dir, "reldir") 

1146 result_submit = [] 

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

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

1149 self.assertEqual( 

1150 set(result_submit), 

1151 { 

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

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

1154 "./tiny_problems.dag", 

1155 "./tiny_problems.dag.rescue001", 

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

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

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

1159 "reldir/001/tiny_problems.node_status", 

1160 }, 

1161 ) 

1162 

1163 def testSubdags(self): 

1164 with temporaryDirectory() as tmp_dir: 

1165 test_tmp_dir = pathlib.Path(tmp_dir) 

1166 submit_dir = test_tmp_dir / "submit" 

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

1168 lssthtc.htc_backup_files(submit_dir) 

1169 result_submit = [] 

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

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

1172 self.assertEqual( 

1173 set(result_submit), 

1174 { 

1175 "./group_failed_1.dag", 

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

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

1178 "./group_failed_1.dag.rescue001", 

1179 "subdags/wms_group_order1_val1a/group_order1_val1a.dag", 

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

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

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

1183 "subdags/wms_group_order1_val1a/group_order1_val1a.node_status", 

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

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

1186 "subdags/wms_group_order1_val1b/group_order1_val1b.dag", 

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

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

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

1190 "subdags/wms_group_order1_val1c/group_order1_val1c.dag", 

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

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

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

1194 "subdags/wms_group_order1_val1c/group_order1_val1c.node_status", 

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

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

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

1198 "001/group_failed_1.info.json", 

1199 "001/group_failed_1.node_status", 

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

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

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

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

1204 }, 

1205 ) 

1206 

1207 

1208class UpdateRescueFileTestCase(unittest.TestCase): 

1209 """Test _update_rescue_file function.""" 

1210 

1211 def testSuccess(self): 

1212 self.maxDiff = None 

1213 with temporaryDirectory() as tmp_dir: 

1214 test_tmp_dir = pathlib.Path(tmp_dir) 

1215 submit_dir = test_tmp_dir / "submit" 

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

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

1218 failed_subdags = lssthtc._update_rescue_file(rescue_file) 

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

1220 with open(rescue_file) as fh: 

1221 lines = fh.readlines() 

1222 results = "".join(lines) 

1223 

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

1225# the u_testuser_DM-46294_group_fail_20250310T160455Z.dag DAG file 

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

1227# Rescue DAG version: 2.0.1 (partial) 

1228# 

1229# Total number of Nodes: 26 

1230# Nodes premarked DONE: 21 

1231# Nodes that failed: 2 

1232# wms_group_order1_val1b,finalJob,<ENDLIST> 

1233 

1234DONE pipetaskInit 

1235DONE label1_val1c_val2a 

1236DONE label1_val1b_val2b 

1237DONE label1_val1b_val2a 

1238DONE label1_val1c_val2b 

1239DONE label1_val1a_val2a 

1240DONE label1_val1a_val2b 

1241DONE label3_val1c_val2a 

1242DONE label3_val1b_val2b 

1243DONE label3_val1b_val2a 

1244DONE label3_val1c_val2b 

1245DONE label3_val1a_val2a 

1246DONE label3_val1a_val2b 

1247DONE wms_group_order1_val1a 

1248DONE label5_val1a_val2a 

1249DONE label5_val1a_val2b 

1250DONE wms_group_order1_val1c 

1251DONE label5_val1c_val2a 

1252DONE label5_val1c_val2b 

1253DONE wms_check_status_wms_group_order1_val1a 

1254DONE wms_check_status_wms_group_order1_val1c 

1255""" 

1256 

1257 self.assertEqual(results, truth) 

1258 

1259 

1260class ReadRescueHeadersTestCase(unittest.TestCase): 

1261 """Test _read_rescue_headers function.""" 

1262 

1263 def testTypical(self): 

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

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

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

1267 

1268 def testEmptyFile(self): 

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

1270 self.assertEqual(result, []) 

1271 

1272 def testOnlyHeaderLines(self): 

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

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

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

1276 

1277 def testFirstLineNotComment(self): 

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

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

1280 self.assertEqual(result, []) 

1281 

1282 def testWhitespaceStripped(self): 

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

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

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

1286 

1287 

1288class WriteRescueHeadersTestCase(unittest.TestCase): 

1289 """Test _write_rescue_headers function.""" 

1290 

1291 def testTypical(self): 

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

1293 outfh = io.StringIO() 

1294 lssthtc._write_rescue_headers(header_lines, outfh) 

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

1296 

1297 def testEmptyList(self): 

1298 outfh = io.StringIO() 

1299 lssthtc._write_rescue_headers([], outfh) 

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

1301 

1302 def testSingleLine(self): 

1303 outfh = io.StringIO() 

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

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

1306 

1307 

1308class UpdateRescueHeadersTestCase(unittest.TestCase): 

1309 """Test _update_rescue_headers function.""" 

1310 

1311 def testWithFailedSubdag(self): 

1312 header_lines = [ 

1313 "# Total number of Nodes: 26", 

1314 "# Nodes premarked DONE: 22", 

1315 "# Nodes that failed: 2", 

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

1317 ] 

1318 result = lssthtc._update_rescue_headers(header_lines) 

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

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

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

1322 

1323 def testNoSubdagFailures(self): 

1324 header_lines = [ 

1325 "# Nodes premarked DONE: 5", 

1326 "# Nodes that failed: 1", 

1327 "# finalJob,<ENDLIST>", 

1328 ] 

1329 result = lssthtc._update_rescue_headers(header_lines) 

1330 self.assertEqual(result, []) 

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

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

1333 

1334 def testMultipleFailedSubdags(self): 

1335 header_lines = [ 

1336 "# Nodes premarked DONE: 10", 

1337 "# Nodes that failed: 3", 

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

1339 ] 

1340 result = lssthtc._update_rescue_headers(header_lines) 

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

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

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

1344 

1345 def testNoFailedNodesLine(self): 

1346 header_lines = [ 

1347 "# Total number of Nodes: 5", 

1348 "# Nodes premarked DONE: 5", 

1349 ] 

1350 original = list(header_lines) 

1351 result = lssthtc._update_rescue_headers(header_lines) 

1352 self.assertEqual(result, []) 

1353 self.assertEqual(header_lines, original) 

1354 

1355 def testEmptyHeader(self): 

1356 result = lssthtc._update_rescue_headers([]) 

1357 self.assertEqual(result, []) 

1358 

1359 

1360class ReadDagStatusTestCase(unittest.TestCase): 

1361 """Test read_dag_status function and read_single_dag_status.""" 

1362 

1363 def testFileMissing(self): 

1364 with temporaryDirectory() as tmp_dir: 

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

1366 _ = lssthtc.read_dag_status(tmp_dir) 

1367 

1368 def testRegular(self): 

1369 with temporaryDirectory() as tmp_dir: 

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

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

1372 results = lssthtc.read_dag_status(submit_dir) 

1373 truth = { 

1374 "JobProcsHeld": 0, 

1375 "NodesPost": 0, 

1376 "JobProcsIdle": 0, 

1377 "NodesTotal": 6, 

1378 "NodesFailed": 2, 

1379 "NodesDone": 3, 

1380 "NodesQueued": 0, 

1381 "NodesPre": 0, 

1382 "NodesFutile": 1, 

1383 "NodesUnready": 0, 

1384 } 

1385 self.assertEqual(results, results | truth) 

1386 

1387 def testSubdags(self): 

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

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

1390 """ 

1391 self.maxDiff = None 

1392 with temporaryDirectory() as tmp_dir: 

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

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

1395 results = lssthtc.read_dag_status(submit_dir) 

1396 truth = { 

1397 "JobProcsHeld": 0, 

1398 "NodesPost": 0, 

1399 "JobProcsIdle": 0, 

1400 "NodesTotal": 34, 

1401 "NodesFailed": 0, 

1402 "NodesDone": 17, 

1403 "NodesQueued": 3, 

1404 "NodesPre": 0, 

1405 "NodesFutile": 0, 

1406 "NodesUnready": 14, 

1407 } 

1408 self.assertEqual(results, results | truth) 

1409 

1410 

1411class ReadDagInfoTestCase(unittest.TestCase): 

1412 """Test read_dag_info function.""" 

1413 

1414 def testFileMissing(self): 

1415 with temporaryDirectory() as tmp_dir: 

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

1417 _ = lssthtc.read_dag_info(tmp_dir) 

1418 

1419 def testSuccess(self): 

1420 with temporaryDirectory() as tmp_dir: 

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

1422 results = lssthtc.read_dag_info(tmp_dir) 

1423 

1424 truth = { 

1425 "test02": { 

1426 "9208.0": { 

1427 "ClusterId": 9208, 

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

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

1430 "bps_project": "dev", 

1431 "bps_payload": "tiny", 

1432 "bps_operator": "testuser", 

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

1434 "bps_provisioning_job": "provisioningJob", 

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

1436 "bps_campaign": "quick", 

1437 "bps_runsite": "testpool", 

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

1439 "bps_run": "u_testuser_tiny_20250213T164427Z", 

1440 "bps_isjob": "True", 

1441 } 

1442 } 

1443 } 

1444 

1445 self.assertEqual(results, truth) 

1446 

1447 def testPermissionError(self): 

1448 with temporaryDirectory() as tmp_dir: 

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

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

1451 mocked_open.side_effect = PermissionError 

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

1453 results = lssthtc.read_dag_info(tmp_dir) 

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

1455 self.assertEqual({}, results) 

1456 

1457 

1458class HtcWriteCondorFileTestCase(unittest.TestCase): 

1459 """Test htc_write_condor_file function.""" 

1460 

1461 def testSuccess(self): 

1462 with temporaryDirectory() as tmp_dir: 

1463 job_name = "job1" 

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

1465 job = { 

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

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

1468 "request_memory": "2000", 

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

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

1471 } 

1472 job_attrs = { 

1473 "bps_job_name": job_name, 

1474 "bps_job_label": "label1", 

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

1476 } 

1477 expected = [ 

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

1479 "arguments=-a -b 2 -c\n", 

1480 "request_memory=2000\n", 

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

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

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

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

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

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

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

1488 "queue\n", 

1489 ] 

1490 

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

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

1493 actual = f.readlines() 

1494 

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

1496 self.assertTrue(filename.exists()) 

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

1498 _ = lssthtc.htc_create_submit_from_file(filename) 

1499 

1500 

1501class HtcCreateSubmitFromDagTestCase(unittest.TestCase): 

1502 """Test htc_create_submit_from_dag function.""" 

1503 

1504 @classmethod 

1505 def setUpClass(cls): 

1506 cls.bindir = None 

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

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

1509 cls.bindir = tempfile.TemporaryDirectory() 

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

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

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

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

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

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

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

1517 

1518 @classmethod 

1519 def tearDownClass(cls): 

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

1521 cls.bindir.cleanup() 

1522 

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

1524 def testMaxIdleEnvVar(self): 

1525 with temporaryDirectory() as tmp_dir: 

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

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

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

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

1530 

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

1532 def testMaxIdleInDAGManConfig(self): 

1533 with temporaryDirectory() as tmp_dir: 

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

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

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

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

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

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

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

1541 

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

1543 def testMaxIdleNotInDAGManConfig(self): 

1544 with temporaryDirectory() as tmp_dir: 

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

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

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

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

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

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

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

1552 

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

1554 def testMaxIdleGiven(self): 

1555 with temporaryDirectory() as tmp_dir: 

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

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

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

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

1560 

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

1562 def testMaxJobsIdleParam(self): 

1563 def _fake_params_contains(key): 

1564 if key == "DAGMAN_MAX_JOBS_IDLE": 

1565 return True 

1566 return False # pragma: no cover 

1567 

1568 def _fake_params_get(key): 

1569 if key == "DAGMAN_MAX_JOBS_IDLE": 

1570 return 16 

1571 return "FAKE_VAL" # pragma: no cover 

1572 

1573 with temporaryDirectory() as tmp_dir: 

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

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

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

1577 mock_param.__contains__.side_effect = _fake_params_contains 

1578 mock_param.__getitem__.side_effect = _fake_params_get 

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

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

1581 

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

1583 def testNoMaxJobsIdle(self): 

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

1585 HTCondor version when no MaxIdle passed to from_dag, not 

1586 checking arguments string here. Instead just making sure 

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

1588 """ 

1589 with temporaryDirectory() as tmp_dir: 

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

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

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

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

1594 mock_param.__contains__.return_value = False 

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

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

1597 

1598 

1599class HtcDagTestCase(unittest.TestCase): 

1600 """Test for HTCDag class.""" 

1601 

1602 def setUp(self): 

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

1604 job.add_job_cmds( 

1605 { 

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

1607 "arguments": "foo", 

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

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

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

1611 } 

1612 ) 

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

1614 

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

1616 self.dag.add_job(job) 

1617 

1618 self.subfile_expected = [ 

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

1620 "arguments=foo\n", 

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

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

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

1624 "queue\n", 

1625 ] 

1626 

1627 def tearDown(self): 

1628 pass 

1629 

1630 def testWriteWithDagConfig(self): 

1631 with temporaryDirectory() as tmp_dir: 

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

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

1634 wms_config_filename = "dagman.conf" 

1635 wms_configurator = dagman_configurator.DagmanConfigurator(config) 

1636 wms_configurator.prepare(wms_config_filename, prefix=tmp_dir) 

1637 wms_configurator.configure(self.dag) 

1638 dagfile_expected = [ 

1639 f"CONFIG {wms_config_filename}\n", 

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

1641 f"DOT {self.dag.name}.dot\n", 

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

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

1644 ] 

1645 

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

1647 

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

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

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

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

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

1653 dagfile_actual = f.readlines() 

1654 self.assertEqual(dagfile_actual, dagfile_expected) 

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

1656 subfile_actual = f.readlines() 

1657 self.assertEqual(subfile_actual, self.subfile_expected) 

1658 

1659 def testWriteWithoutDagConfig(self): 

1660 with temporaryDirectory() as tmp_dir: 

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

1662 dagfile_expected = [ 

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

1664 f"DOT {self.dag.name}.dot\n", 

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

1666 ] 

1667 

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

1669 

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

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

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

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

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

1675 dagfile_actual = f.readlines() 

1676 self.assertEqual(dagfile_actual, dagfile_expected) 

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

1678 subfile_actual = f.readlines() 

1679 self.assertEqual(subfile_actual, self.subfile_expected) 

1680 

1681 

1682if __name__ == "__main__": 

1683 unittest.main()