Coverage for tests/test_lssthtc.py: 99%
829 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-15 09:04 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-15 09:04 +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."""
29import io
30import logging
31import os
32import pathlib
33import stat
34import tempfile
35import unittest
36from shutil import copy2, copytree, ignore_patterns, rmtree, which
38import htcondor
39from dag_test_utils import make_lazy_dag
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
47logger = logging.getLogger("lsst.ctrl.bps.htcondor")
48TESTDIR = os.path.abspath(os.path.dirname(__file__))
51class TestLsstHtc(unittest.TestCase):
52 """Test basic usage."""
54 def testHtcEscapeInt(self):
55 self.assertEqual(lssthtc.htc_escape(100), 100)
57 def testHtcEscapeDouble(self):
58 self.assertEqual(lssthtc.htc_escape('"double"'), '""double""')
60 def testHtcEscapeSingle(self):
61 self.assertEqual(lssthtc.htc_escape("'single'"), "''single''")
63 def testHtcEscapeNoSideEffect(self):
64 val = "'val'"
65 self.assertEqual(lssthtc.htc_escape(val), "''val''")
66 self.assertEqual(val, "'val'")
68 def testHtcEscapeQuot(self):
69 self.assertEqual(lssthtc.htc_escape(""val""), '"val"')
71 def testHtcVersion(self):
72 ver = lssthtc.htc_version()
73 self.assertRegex(ver, r"^\d+\.\d+\.\d+$")
76class HtcTweakJobInfoTestCase(unittest.TestCase):
77 """Test the function responsible for massaging job information."""
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 }
91 def tearDown(self):
92 self.log_dir.cleanup()
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())
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)
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)
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)
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)
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)
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)
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)
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)
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)
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)
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])
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])
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'")
189class HtcCheckDagmanOutputTestCase(unittest.TestCase):
190 """Test htc_check_dagman_output function."""
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)
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)
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)
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)
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)
225class SummarizeDagTestCase(unittest.TestCase):
226 """Test summarize_dag function."""
228 def test_no_dag_file(self):
229 with temporaryDirectory() as tmp_dir:
230 summary, job_name_to_pipetask, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
231 self.assertFalse(len(job_name_to_pipetask))
232 self.assertFalse(len(job_name_to_type))
233 self.assertFalse(summary)
235 def test_success(self):
236 with temporaryDirectory() as tmp_dir:
237 copy2(f"{TESTDIR}/data/good.dag", tmp_dir)
238 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
239 self.assertEqual(summary, "pipetaskInit:1;label1:1;label2:1;label3:1;finalJob:1")
240 self.assertEqual(
241 job_name_to_label,
242 {
243 "pipetaskInit": "pipetaskInit",
244 "0682f8f9-12f0-40a5-971e-8b30c7231e5c_label1_val1_val2": "label1",
245 "d0305e2d-f164-4a85-bd24-06afe6c84ed9_label2_val1_val2": "label2",
246 "2806ecc9-1bba-4362-8fff-ab4e6abb9f83_label3_val1_val2": "label3",
247 "finalJob": "finalJob",
248 },
249 )
250 self.assertEqual(
251 job_name_to_type,
252 {
253 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
254 "0682f8f9-12f0-40a5-971e-8b30c7231e5c_label1_val1_val2": lssthtc.WmsNodeType.PAYLOAD,
255 "d0305e2d-f164-4a85-bd24-06afe6c84ed9_label2_val1_val2": lssthtc.WmsNodeType.PAYLOAD,
256 "2806ecc9-1bba-4362-8fff-ab4e6abb9f83_label3_val1_val2": lssthtc.WmsNodeType.PAYLOAD,
257 "finalJob": lssthtc.WmsNodeType.FINAL,
258 },
259 )
261 def test_service(self):
262 with temporaryDirectory() as tmp_dir:
263 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag", tmp_dir)
264 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
265 self.assertEqual(summary, "pipetaskInit:1;label1:2;label2:2;finalJob:1")
266 self.assertEqual(
267 job_name_to_label,
268 {
269 "pipetaskInit": "pipetaskInit",
270 "057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b": "label1",
271 "4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a": "label1",
272 "40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b": "label2",
273 "696ee50d-e711-40d6-9caf-ee29ae4a656d_label2_val1a_val2a": "label2",
274 "finalJob": "finalJob",
275 "provisioningJob": "provisioningJob",
276 },
277 )
278 self.assertEqual(
279 job_name_to_type,
280 {
281 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
282 "057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
283 "4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
284 "40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
285 "696ee50d-e711-40d6-9caf-ee29ae4a656d_label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
286 "finalJob": lssthtc.WmsNodeType.FINAL,
287 "provisioningJob": lssthtc.WmsNodeType.SERVICE,
288 },
289 )
291 def test_noop(self):
292 with temporaryDirectory() as tmp_dir:
293 copy2(f"{TESTDIR}/data/noop_running_1/noop_running_1.dag", tmp_dir)
294 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
295 self.assertEqual(
296 set(summary.split(";")),
297 {"pipetaskInit:1", "label1:6", "label2:6", "label3:6", "label4:6", "label5:6", "finalJob:1"},
298 )
299 self.assertEqual(
300 job_name_to_label,
301 {
302 "label1_val1a_val2a": "label1",
303 "label1_val1a_val2b": "label1",
304 "label1_val1b_val2a": "label1",
305 "label1_val1b_val2b": "label1",
306 "label1_val1c_val2a": "label1",
307 "label1_val1c_val2b": "label1",
308 "label2_val1a_val2a": "label2",
309 "label2_val1a_val2b": "label2",
310 "label2_val1b_val2a": "label2",
311 "label2_val1b_val2b": "label2",
312 "label2_val1c_val2a": "label2",
313 "label2_val1c_val2b": "label2",
314 "label3_val1a_val2a": "label3",
315 "label3_val1a_val2b": "label3",
316 "label3_val1b_val2a": "label3",
317 "label3_val1b_val2b": "label3",
318 "label3_val1c_val2a": "label3",
319 "label3_val1c_val2b": "label3",
320 "label4_val1a_val2a": "label4",
321 "label4_val1a_val2b": "label4",
322 "label4_val1b_val2a": "label4",
323 "label4_val1b_val2b": "label4",
324 "label4_val1c_val2a": "label4",
325 "label4_val1c_val2b": "label4",
326 "label5_val1a_val2a": "label5",
327 "label5_val1a_val2b": "label5",
328 "label5_val1b_val2a": "label5",
329 "label5_val1b_val2b": "label5",
330 "label5_val1c_val2a": "label5",
331 "label5_val1c_val2b": "label5",
332 "finalJob": "finalJob",
333 "pipetaskInit": "pipetaskInit",
334 "wms_noop_order1_val1a": "order1",
335 "wms_noop_order1_val1b": "order1",
336 },
337 )
338 self.assertEqual(
339 job_name_to_type,
340 {
341 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
342 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
343 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
344 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
345 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
346 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
347 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
348 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
349 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
350 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
351 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
352 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
353 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
354 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
355 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
356 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
357 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
358 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
359 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
360 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
361 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
362 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
363 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
364 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
365 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
366 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
367 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
368 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
369 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
370 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
371 "finalJob": lssthtc.WmsNodeType.FINAL,
372 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
373 "wms_noop_order1_val1a": lssthtc.WmsNodeType.NOOP,
374 "wms_noop_order1_val1b": lssthtc.WmsNodeType.NOOP,
375 },
376 )
378 def test_subdags(self):
379 self.maxDiff = None
380 with temporaryDirectory() as tmp_dir:
381 submit_dir = os.path.join(tmp_dir, "group_running_1")
382 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
383 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(submit_dir)
384 self.assertEqual(
385 set(summary.split(";")),
386 {"pipetaskInit:1", "label1:6", "label2:6", "label3:6", "label4:6", "label5:6", "finalJob:1"},
387 )
389 self.assertEqual(
390 job_name_to_label,
391 {
392 "pipetaskInit": "pipetaskInit",
393 "label1_val1b_val2a": "label1",
394 "label1_val1c_val2a": "label1",
395 "label1_val1a_val2b": "label1",
396 "label1_val1b_val2b": "label1",
397 "label1_val1c_val2b": "label1",
398 "label1_val1a_val2a": "label1",
399 "label2_val1a_val2b": "label2",
400 "label2_val1a_val2a": "label2",
401 "label2_val1b_val2a": "label2",
402 "label2_val1b_val2b": "label2",
403 "label2_val1c_val2a": "label2",
404 "label2_val1c_val2b": "label2",
405 "label3_val1b_val2a": "label3",
406 "label3_val1c_val2a": "label3",
407 "label3_val1a_val2b": "label3",
408 "label3_val1b_val2b": "label3",
409 "label3_val1c_val2b": "label3",
410 "label3_val1a_val2a": "label3",
411 "label4_val1a_val2b": "label4",
412 "label4_val1a_val2a": "label4",
413 "label4_val1b_val2a": "label4",
414 "label4_val1b_val2b": "label4",
415 "label4_val1c_val2a": "label4",
416 "label4_val1c_val2b": "label4",
417 "label5_val1a_val2b": "label5",
418 "label5_val1a_val2a": "label5",
419 "label5_val1b_val2a": "label5",
420 "label5_val1b_val2b": "label5",
421 "label5_val1c_val2a": "label5",
422 "label5_val1c_val2b": "label5",
423 "finalJob": "finalJob",
424 "provisioningJob": "provisioningJob",
425 "wms_group_order1_val1a": "order1",
426 "wms_group_order1_val1b": "order1",
427 "wms_group_order1_val1c": "order1",
428 "wms_check_status_wms_group_order1_val1a": "order1",
429 "wms_check_status_wms_group_order1_val1b": "order1",
430 "wms_check_status_wms_group_order1_val1c": "order1",
431 },
432 )
434 self.assertEqual(
435 job_name_to_type,
436 {
437 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
438 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
439 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
440 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
441 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
442 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
443 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
444 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
445 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
446 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
447 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
448 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
449 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
450 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
451 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
452 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
453 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
454 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
455 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
456 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
457 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
458 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
459 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
460 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
461 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
462 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
463 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
464 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
465 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
466 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
467 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
468 "finalJob": lssthtc.WmsNodeType.FINAL,
469 "provisioningJob": lssthtc.WmsNodeType.SERVICE,
470 "wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG,
471 "wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG,
472 "wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG,
473 "wms_check_status_wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG_CHECK,
474 "wms_check_status_wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG_CHECK,
475 "wms_check_status_wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG_CHECK,
476 },
477 )
480class ReadDagNodesLogTestCase(unittest.TestCase):
481 """Test read_dag_nodes_log function."""
483 def setUp(self):
484 self.tmpdir = tempfile.mkdtemp()
486 def tearDown(self):
487 rmtree(self.tmpdir, ignore_errors=True)
489 def testFileMissing(self):
490 with self.assertRaisesRegex(FileNotFoundError, "DAGMan node log not found in"):
491 _ = lssthtc.read_dag_nodes_log(self.tmpdir)
493 def testRegular(self):
494 with temporaryDirectory() as tmp_dir:
495 submit_dir = os.path.join(tmp_dir, "tiny_problems")
496 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
497 results = lssthtc.read_dag_nodes_log(submit_dir)
498 self.assertEqual(results["9231.0"]["Cluster"], 9231)
499 self.assertEqual(results["9231.0"]["Proc"], 0)
500 self.assertEqual(results["9231.0"]["ToE"]["ExitCode"], 1)
501 self.assertEqual(len(results), 6)
503 def testSubdags(self):
504 """Making sure it gets data from subdag dirs and doesn't
505 fail if some subdags haven't started running yet.
506 """
507 with temporaryDirectory() as tmp_dir:
508 submit_dir = os.path.join(tmp_dir, "group_running_1")
509 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
510 results = lssthtc.read_dag_nodes_log(submit_dir)
511 # main dag
512 self.assertEqual(results["10094.0"]["Cluster"], 10094)
513 # subdag
514 self.assertEqual(results["10112.0"]["Cluster"], 10112)
515 self.assertEqual(results["10116.0"]["Cluster"], 10116)
518class ReadNodeStatusTestCase(unittest.TestCase):
519 """Test read_node_status function."""
521 def setUp(self):
522 self.tmpdir = tempfile.mkdtemp()
524 def tearDown(self):
525 rmtree(self.tmpdir, ignore_errors=True)
527 def testServiceJobNotSubmitted(self):
528 # tiny_prov_no_submit files have successful workflow
529 # but provisioningJob could not submit.
530 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag.nodes.log", self.tmpdir)
531 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag.dagman.log", self.tmpdir)
532 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.node_status", self.tmpdir)
533 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag", self.tmpdir)
535 jobs = lssthtc.read_node_status(self.tmpdir)
536 found = [
537 id_
538 for id_ in jobs
539 if jobs[id_].get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN) == lssthtc.WmsNodeType.SERVICE
540 ]
541 self.assertEqual(len(found), 1)
542 self.assertEqual(jobs[found[0]]["DAGNodeName"], "provisioningJob")
543 self.assertEqual(jobs[found[0]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
545 def testMissingStatusFile(self):
546 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag.nodes.log", self.tmpdir)
547 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag.dagman.log", self.tmpdir)
548 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag", self.tmpdir)
550 jobs = lssthtc.read_node_status(self.tmpdir)
551 self.assertEqual(len(jobs), 7)
552 self.assertEqual(jobs["9230.0"]["DAGNodeName"], "pipetaskInit")
553 self.assertEqual(jobs["9230.0"]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD)
554 found = [
555 id_
556 for id_ in jobs
557 if jobs[id_].get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN) == lssthtc.WmsNodeType.SERVICE
558 ]
559 self.assertEqual(len(found), 1)
560 self.assertEqual(jobs[found[0]]["DAGNodeName"], "provisioningJob")
562 def testSubdagsRunning(self):
563 with temporaryDirectory() as tmp_dir:
564 test_tmp_dir = pathlib.Path(tmp_dir)
565 submit_dir = test_tmp_dir / "submit"
566 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
567 jobs = lssthtc.read_node_status(submit_dir)
568 self.assertEqual(len(jobs), 39) # includes non-payload jobs
569 # not guaranteed ids are same, so use names instead
570 job_name_to_id = {}
571 for id_, info in jobs.items():
572 job_name_to_id[info.get("DAGNodeName", id_)] = id_
573 job_type_to_names = {}
574 for id_, info in jobs.items():
575 job_type_to_names.setdefault(
576 info.get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN), set()
577 ).add(info.get("DAGNodeName", id_))
579 # check counts
580 self.assertNotIn(lssthtc.WmsNodeType.NOOP, job_type_to_names)
581 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.PAYLOAD]), 31)
582 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.FINAL]), 1)
583 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SERVICE]), 1)
584 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG]), 3)
585 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG_CHECK]), 3)
587 # spot check some statuses
588 self.assertEqual(
589 jobs[job_name_to_id["label3_val1a_val2b"]]["NodeStatus"], lssthtc.NodeStatus.DONE
590 )
591 self.assertEqual(
592 jobs[job_name_to_id["wms_group_order1_val1a"]]["NodeStatus"], lssthtc.NodeStatus.SUBMITTED
593 )
594 self.assertEqual(
595 jobs[job_name_to_id["label5_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY
596 )
597 self.assertEqual(
598 jobs[job_name_to_id["label2_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.DONE
599 )
601 def testSubdagsFailed(self):
602 with temporaryDirectory() as tmp_dir:
603 test_tmp_dir = pathlib.Path(tmp_dir)
604 submit_dir = test_tmp_dir / "submit"
605 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
606 jobs = lssthtc.read_node_status(submit_dir)
607 self.assertEqual(len(jobs), 39)
608 # not guaranteed ids are same, so use names instead
609 job_name_to_id = {}
610 for id_, info in jobs.items():
611 job_name_to_id[info.get("DAGNodeName", id_)] = id_
612 job_type_to_names = {}
613 for id_, info in jobs.items():
614 job_type_to_names.setdefault(
615 info.get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN), set()
616 ).add(info.get("DAGNodeName", id_))
618 # check counts
619 self.assertNotIn(lssthtc.WmsNodeType.NOOP, job_type_to_names)
620 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.PAYLOAD]), 31)
621 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.FINAL]), 1)
622 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SERVICE]), 1)
623 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG]), 3)
624 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG_CHECK]), 3)
626 # spot check some statuses
627 self.assertEqual(
628 jobs[job_name_to_id["label3_val1a_val2b"]]["NodeStatus"], lssthtc.NodeStatus.DONE
629 )
630 self.assertEqual(
631 jobs[job_name_to_id["wms_group_order1_val1a"]]["NodeStatus"], lssthtc.NodeStatus.DONE
632 )
633 self.assertEqual(
634 jobs[job_name_to_id["label5_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.DONE
635 )
637 self.assertEqual(
638 jobs[job_name_to_id["label5_val1b_val2a"]]["NodeStatus"], lssthtc.NodeStatus.FUTILE
639 )
640 self.assertEqual(
641 jobs[job_name_to_id["wms_group_order1_val1b"]]["NodeStatus"], lssthtc.NodeStatus.DONE
642 )
643 self.assertEqual(
644 jobs[job_name_to_id["wms_check_status_wms_group_order1_val1b"]]["NodeStatus"],
645 lssthtc.NodeStatus.ERROR,
646 )
649class ReadSingleNodeStatusTestCase(unittest.TestCase):
650 """Test read_single_node_status function."""
652 def setUp(self):
653 self.tmpdir = tempfile.mkdtemp()
655 def tearDown(self):
656 rmtree(self.tmpdir, ignore_errors=True)
658 def _copyFiles(self, data_subdir, suffixes):
659 """Copy files with given suffixes from tests/data/<data_subdir>/."""
660 for suffix in suffixes:
661 copy2(f"{TESTDIR}/data/{data_subdir}/{data_subdir}{suffix}", self.tmpdir)
663 def _jobNameToId(self, jobs):
664 return {info["DAGNodeName"]: id_ for id_, info in jobs.items()}
666 def testAllDone(self):
667 self._copyFiles(
668 "tiny_success",
669 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
670 )
671 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
672 jobs = lssthtc.read_single_node_status(filename, -1)
674 self.assertEqual(len(jobs), 5)
675 name_to_id = self._jobNameToId(jobs)
677 # All four submitted nodes are marked DONE.
678 for name in [
679 "pipetaskInit",
680 "5bba27bd-8df7-4668-a9c5-e911192c5cdb_label1_val1_val2",
681 "0b225f1f-6edf-4380-b546-76c97947a88f_label2_val1_val2",
682 "finalJob",
683 ]:
684 self.assertIn(name, name_to_id, msg=f"Missing job {name}")
685 self.assertEqual(
686 jobs[name_to_id[name]]["NodeStatus"],
687 lssthtc.NodeStatus.DONE,
688 msg=f"Expected DONE for {name}",
689 )
691 # Service job not tracked by node_status; it came from the event log so
692 # it has a real positive ClusterId but no NodeStatus field.
693 self.assertIn("provisioningJob", name_to_id)
694 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
696 # Spot-check labels and types.
697 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["bps_job_label"], "pipetaskInit")
698 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD)
699 self.assertEqual(jobs[name_to_id["finalJob"]]["wms_node_type"], lssthtc.WmsNodeType.FINAL)
700 self.assertEqual(jobs[name_to_id["provisioningJob"]]["wms_node_type"], lssthtc.WmsNodeType.SERVICE)
702 # DAGManJobID is populated from the dagman log for every job.
703 for job in jobs.values():
704 self.assertIn("DAGManJobID", job)
706 def testMixedStatuses(self):
707 self._copyFiles(
708 "tiny_problems",
709 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
710 )
711 filename = pathlib.Path(self.tmpdir) / "tiny_problems.node_status"
712 jobs = lssthtc.read_single_node_status(filename, -1)
714 self.assertEqual(len(jobs), 7)
715 name_to_id = self._jobNameToId(jobs)
717 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
718 self.assertEqual(
719 jobs[name_to_id["057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b"]]["NodeStatus"],
720 lssthtc.NodeStatus.ERROR,
721 )
722 self.assertEqual(
723 jobs[name_to_id["4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a"]]["NodeStatus"],
724 lssthtc.NodeStatus.DONE,
725 )
726 self.assertEqual(
727 jobs[name_to_id["40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b"]]["NodeStatus"],
728 lssthtc.NodeStatus.FUTILE,
729 )
730 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.ERROR)
731 # Service job not tracked by node_status; came from event log so
732 # it has a real positive ClusterId but no NodeStatus field.
733 self.assertIn("provisioningJob", name_to_id)
734 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
736 def testRunningWorkflow(self):
737 self._copyFiles(
738 "tiny_running",
739 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
740 )
741 filename = pathlib.Path(self.tmpdir) / "tiny_running.node_status"
742 jobs = lssthtc.read_single_node_status(filename, -1)
744 self.assertEqual(len(jobs), 5)
745 name_to_id = self._jobNameToId(jobs)
747 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
748 self.assertEqual(
749 jobs[name_to_id["ca27ea57-c014-44c1-838a-78c06bc3ec1b_label1_val1_val2"]]["NodeStatus"],
750 lssthtc.NodeStatus.SUBMITTED,
751 )
752 self.assertEqual(
753 jobs[name_to_id["dbf919fa-5453-4b05-8806-ad6390fda0a3_label2_val1_val2"]]["NodeStatus"],
754 lssthtc.NodeStatus.NOT_READY,
755 )
756 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
757 # Service job appeared in the event log; has a real positive ClusterId.
758 self.assertIn("provisioningJob", name_to_id)
759 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
761 def testMissingNodeStatusFile(self):
762 # Omit the .node_status file; jobs must be built from the event log
763 # and dag.
764 self._copyFiles(
765 "tiny_problems",
766 [".dag", ".dag.dagman.log", ".dag.nodes.log"],
767 )
768 filename = pathlib.Path(self.tmpdir) / "tiny_problems.node_status"
769 jobs = lssthtc.read_single_node_status(filename, -1)
771 self.assertEqual(len(jobs), 7)
772 name_to_id = self._jobNameToId(jobs)
774 # Jobs that appeared in the event log have real (positive) cluster IDs.
775 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["DAGNodeName"], "pipetaskInit")
776 self.assertGreater(jobs[name_to_id["pipetaskInit"]]["ClusterId"], 0)
778 # The service job appeared in the event log and has a real positive ID.
779 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
781 # All jobs carry the correct label and type from the dag file.
782 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD)
783 self.assertEqual(jobs[name_to_id["provisioningJob"]]["wms_node_type"], lssthtc.WmsNodeType.SERVICE)
785 def testMissingLogFiles(self):
786 # Omit both log files; every job should get a fake negative ID.
787 self._copyFiles("tiny_success", [".dag", ".node_status"])
788 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
789 jobs = lssthtc.read_single_node_status(filename, -1)
791 self.assertEqual(len(jobs), 5)
792 for job in jobs.values():
793 self.assertLess(job["ClusterId"], 0)
795 # NodeStatus values from the node_status file must still be preserved.
796 name_to_id = self._jobNameToId(jobs)
797 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
798 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
799 self.assertEqual(jobs[name_to_id["provisioningJob"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
801 def testInitFakeId(self):
802 # Verify fake IDs count down from the given starting value.
803 self._copyFiles("tiny_success", [".dag", ".node_status"])
804 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
805 init_fake_id = -10
806 jobs = lssthtc.read_single_node_status(filename, init_fake_id)
808 self.assertEqual(len(jobs), 5)
809 cluster_ids = [job["ClusterId"] for job in jobs.values()]
810 # All IDs must be at most init_fake_id (i.e., -10 or lower).
811 for cid in cluster_ids:
812 self.assertLessEqual(cid, init_fake_id)
813 # All IDs must be unique.
814 self.assertEqual(len(set(cluster_ids)), len(cluster_ids))
816 def testFromDagJobAttribute(self):
817 self._copyFiles(
818 "tiny_success",
819 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
820 )
821 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
822 jobs = lssthtc.read_single_node_status(filename, -1)
824 for id_, job in jobs.items():
825 self.assertEqual(
826 job["from_dag_job"],
827 "wms_tiny_success",
828 msg=f"Job {id_} has wrong from_dag_job",
829 )
831 def testServiceJobPlaceholder(self):
832 self._copyFiles(
833 "tiny_prov_no_submit",
834 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
835 )
836 filename = pathlib.Path(self.tmpdir) / "tiny_prov_no_submit.node_status"
837 jobs = lssthtc.read_single_node_status(filename, -1)
839 service_jobs = [
840 (id_, info)
841 for id_, info in jobs.items()
842 if info.get("wms_node_type") == lssthtc.WmsNodeType.SERVICE
843 ]
844 self.assertEqual(len(service_jobs), 1)
845 service_id, service_job = service_jobs[0]
846 self.assertEqual(service_job["DAGNodeName"], "provisioningJob")
847 self.assertEqual(service_job["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
848 self.assertLess(service_job["ClusterId"], 0)
851class HTCJobTestCase(unittest.TestCase):
852 """Test HTCJob methods."""
854 def testWriteDagCommandsPayload(self):
855 job = lssthtc.HTCJob(
856 "job1",
857 "label1",
858 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
859 {"dir": "jobs/label1"},
860 )
861 job.subfile = "job1.sub"
863 mockfh = io.StringIO()
864 job.write_dag_commands(mockfh, "../..")
865 self.assertIn('JOB job1 "job1.sub" DIR "../../jobs/label1"', mockfh.getvalue())
867 def testWriteDagCommandsNotJob(self):
868 # Testing giving command_name, no dag_rel_path and no dir
869 job = lssthtc.HTCJob(
870 "finalJob",
871 "finalJob",
872 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
873 )
874 job.subfile = "jobs/finalJob/finalJob.sub"
875 mockfh = io.StringIO()
876 job.write_dag_commands(mockfh, "", "FINAL")
877 self.assertIn('FINAL finalJob "jobs/finalJob/finalJob.sub"', mockfh.getvalue())
879 def testWriteDagCommandsNoop(self):
880 job = lssthtc.HTCJob("wms_noop_job1", "label1", {}, {"noop": True})
881 job.subfile = "notthere.sub"
882 mockfh = io.StringIO()
883 job.write_dag_commands(mockfh, "")
884 self.assertIn("NOOP", mockfh.getvalue())
886 def testWriteSubmitFile(self):
887 job = lssthtc.HTCJob(
888 "job1",
889 "label1",
890 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
891 )
892 with temporaryDirectory() as tmp_dir:
893 filename = pathlib.Path(tmp_dir) / "label1/job1.sub"
894 job.write_submit_file(filename.parent)
895 self.assertTrue(filename.exists())
896 # Try to make Submit object from file to find any syntax issues
897 _ = lssthtc.htc_create_submit_from_file(filename)
899 def testWriteSubmitFileExists(self):
900 job = lssthtc.HTCJob(
901 "job1",
902 "label1",
903 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
904 )
905 with temporaryDirectory() as tmp_dir:
906 filename = pathlib.Path(tmp_dir) / "job1.sub"
907 job.subfile = filename
908 with open(filename, "w"):
909 pass # make empty file
910 job.write_submit_file(filename.parent)
911 # make sure didn't overwrite file
912 self.assertEqual(filename.stat().st_size, 0, "Incorrectly overwrote existing file")
915class HtcWriteJobCommands(unittest.TestCase):
916 """Test _htc_write_job_commands function."""
918 def testAllCommands(self):
919 dag_cmds = {
920 "pre": {
921 "defer": {"status": 1, "time": 120},
922 "debug": {"filename": "debug_pre.txt", "type": "ALL"},
923 "executable": "exec1",
924 "arguments": "arg1 arg2",
925 },
926 "post": {
927 "defer": {"status": 2, "time": 180},
928 "debug": {"filename": "debug_post.txt", "type": "ALL"},
929 "executable": "exec2",
930 "arguments": "arg3 arg4",
931 },
932 "vars": {"num": 8, "spaces": "a space"},
933 "pre_skip": "1",
934 "retry": 3,
935 "retry_unless_exit": 1,
936 "abort_dag_on": {"node_exit": 100, "abort_exit": 4},
937 "priority": 123,
938 }
940 truth = """SCRIPT DEFER 1 120 DEBUG debug_pre.txt ALL PRE job1 exec1 arg1 arg2
941SCRIPT DEFER 2 180 DEBUG debug_post.txt ALL POST job1 exec2 arg3 arg4
942VARS job1 num="8"
943VARS job1 spaces="a space"
944PRE_SKIP job1 1
945RETRY job1 3 UNLESS-EXIT 1
946ABORT-DAG-ON job1 100 RETURN 4
947PRIORITY job1 123
948"""
949 mockfh = io.StringIO()
950 lssthtc._htc_write_job_commands(mockfh, "job1", dag_cmds)
951 self.assertEqual(mockfh.getvalue(), truth)
953 def testPartialCommands(self):
954 # Trigger skipping the inner if clauses.
955 dag_cmds = {
956 "pre": {
957 "executable": "exec1",
958 },
959 "post": {
960 "executable": "exec2",
961 },
962 "vars": {"num": 8, "spaces": "a space"},
963 "pre_skip": "1",
964 "retry": 3,
965 }
967 truth = """SCRIPT PRE job1 exec1
968SCRIPT POST job1 exec2
969VARS job1 num="8"
970VARS job1 spaces="a space"
971PRE_SKIP job1 1
972RETRY job1 3
973"""
974 mockfh = io.StringIO()
975 lssthtc._htc_write_job_commands(mockfh, "job1", dag_cmds)
976 self.assertEqual(mockfh.getvalue(), truth)
978 def testNoCommands(self):
979 dag_cmds = {}
980 mockfh = io.StringIO()
981 lssthtc._htc_write_job_commands(mockfh, "job2", dag_cmds)
982 self.assertEqual(mockfh.getvalue(), "")
984 def testFinal(self):
985 self.maxDiff = None
986 dag_cmds = {
987 "pre": {
988 "defer": {"status": 1, "time": 120},
989 "debug": {"filename": "debug_pre.txt", "type": "ALL"},
990 "executable": "exec1",
991 "arguments": "arg1 arg2",
992 },
993 "post": {
994 "defer": {"status": 2, "time": 180},
995 "debug": {"filename": "debug_post.txt", "type": "ALL"},
996 "executable": "exec2",
997 "arguments": "arg3 arg4",
998 },
999 "vars": {"num": 8, "spaces": "a space"},
1000 "pre_skip": "1",
1001 "retry": 3,
1002 "retry_unless_exit": 1,
1003 "abort_dag_on": {"node_exit": 100, "abort_exit": 4},
1004 "priority": 123,
1005 }
1007 truth = """SCRIPT DEFER 1 120 DEBUG debug_pre.txt ALL PRE finalJob exec1 arg1 arg2
1008SCRIPT DEFER 2 180 DEBUG debug_post.txt ALL POST finalJob exec2 arg3 arg4
1009VARS finalJob num="8"
1010VARS finalJob spaces="a space"
1011PRE_SKIP finalJob 1
1012"""
1013 mockfh = io.StringIO()
1014 lssthtc._htc_write_job_commands(mockfh, "finalJob", dag_cmds, "FINAL")
1015 self.assertEqual(mockfh.getvalue(), truth)
1018class HTCBackupFilesSinglePathTestCase(unittest.TestCase):
1019 """Test htc_backup_files_single_path function."""
1021 def testSrcDestSame(self):
1022 with temporaryDirectory() as tmp_dir:
1023 with self.assertRaisesRegex(
1024 RuntimeError, "Destination directory is same as the source directory"
1025 ):
1026 lssthtc.htc_backup_files_single_path(tmp_dir, tmp_dir)
1028 def testSuccess(self):
1029 with temporaryDirectory() as tmp_dir:
1030 test_tmp_dir = pathlib.Path(tmp_dir)
1031 submit_dir = test_tmp_dir / "the_src_dir"
1032 copytree(f"{TESTDIR}/data/tiny_success", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1033 backup_dir = test_tmp_dir / "the_dest_dir"
1034 backup_dir.mkdir()
1035 lssthtc.htc_backup_files_single_path(submit_dir, backup_dir)
1036 result_submit = []
1037 for root, _, files in os.walk(submit_dir):
1038 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1039 self.assertEqual(
1040 set(result_submit),
1041 {
1042 "./tiny_success.dag.dagman.log",
1043 "./tiny_success.dag.dagman.out",
1044 "./tiny_success.dag",
1045 },
1046 )
1047 result_backup = []
1048 for root, _, files in os.walk(backup_dir):
1049 result_backup.extend([str(os.path.join(os.path.relpath(root, backup_dir), f)) for f in files])
1050 self.assertEqual(
1051 set(result_backup),
1052 {
1053 "./tiny_success.info.json",
1054 "./tiny_success.dag.metrics",
1055 "./tiny_success.dag.nodes.log",
1056 "./tiny_success.node_status",
1057 },
1058 )
1061class HTCBackupFilesTestCase(unittest.TestCase):
1062 """Test htc_backup_files function."""
1064 def testDirectoryNotFound(self):
1065 with temporaryDirectory() as tmp_dir:
1066 test_tmp_dir = pathlib.Path(tmp_dir)
1067 submit_dir = test_tmp_dir / "submit"
1068 with self.assertRaises(FileNotFoundError):
1069 lssthtc.htc_backup_files(submit_dir)
1071 def testSuccess(self):
1072 with temporaryDirectory() as tmp_dir:
1073 test_tmp_dir = pathlib.Path(tmp_dir)
1074 submit_dir = test_tmp_dir / "submit"
1075 copytree(f"{TESTDIR}/data/tiny_success", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1076 lssthtc.htc_backup_files(submit_dir)
1077 result_submit = []
1078 for root, _, files in os.walk(submit_dir):
1079 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1080 self.assertEqual(
1081 set(result_submit),
1082 {
1083 "./tiny_success.dag.dagman.log",
1084 "./tiny_success.dag.dagman.out",
1085 "./tiny_success.dag",
1086 "000/tiny_success.info.json",
1087 "000/tiny_success.dag.metrics",
1088 "000/tiny_success.dag.nodes.log",
1089 "000/tiny_success.node_status",
1090 },
1091 )
1093 def testDestNotInSubmitDir(self):
1094 with temporaryDirectory() as tmp_dir:
1095 test_tmp_dir = pathlib.Path(tmp_dir)
1096 submit_dir = test_tmp_dir / "submit"
1097 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1098 with self.assertLogs("lsst.ctrl.bps.htcondor", level="WARNING") as cm:
1099 lssthtc.htc_backup_files(submit_dir, test_tmp_dir / "backup")
1100 self.assertIn("Invalid backup location:", cm.output[-1])
1101 lssthtc.htc_backup_files(submit_dir)
1102 result_submit = []
1103 for root, _, files in os.walk(submit_dir):
1104 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1105 self.assertEqual(
1106 set(result_submit),
1107 {
1108 "./tiny_problems.dag.dagman.log",
1109 "./tiny_problems.dag.dagman.out",
1110 "./tiny_problems.dag",
1111 "./tiny_problems.dag.rescue001",
1112 "001/tiny_problems.info.json",
1113 "001/tiny_problems.dag.metrics",
1114 "001/tiny_problems.dag.nodes.log",
1115 "001/tiny_problems.node_status",
1116 },
1117 )
1119 def testDestInSubmitDir(self):
1120 with temporaryDirectory() as tmp_dir:
1121 test_tmp_dir = pathlib.Path(tmp_dir)
1122 submit_dir = test_tmp_dir / "submit"
1123 backup_dir = submit_dir / "subdir"
1124 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1125 lssthtc.htc_backup_files(submit_dir, backup_dir)
1126 result_submit = []
1127 for root, _, files in os.walk(submit_dir):
1128 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1129 self.assertEqual(
1130 set(result_submit),
1131 {
1132 "./tiny_problems.dag.dagman.log",
1133 "./tiny_problems.dag.dagman.out",
1134 "./tiny_problems.dag",
1135 "./tiny_problems.dag.rescue001",
1136 "subdir/001/tiny_problems.info.json",
1137 "subdir/001/tiny_problems.dag.metrics",
1138 "subdir/001/tiny_problems.dag.nodes.log",
1139 "subdir/001/tiny_problems.node_status",
1140 },
1141 )
1143 def testRelativeSubdir(self):
1144 with temporaryDirectory() as tmp_dir:
1145 test_tmp_dir = pathlib.Path(tmp_dir)
1146 submit_dir = test_tmp_dir / "submit"
1147 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1148 lssthtc.htc_backup_files(submit_dir, "reldir")
1149 result_submit = []
1150 for root, _, files in os.walk(submit_dir):
1151 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1152 self.assertEqual(
1153 set(result_submit),
1154 {
1155 "./tiny_problems.dag.dagman.log",
1156 "./tiny_problems.dag.dagman.out",
1157 "./tiny_problems.dag",
1158 "./tiny_problems.dag.rescue001",
1159 "reldir/001/tiny_problems.info.json",
1160 "reldir/001/tiny_problems.dag.metrics",
1161 "reldir/001/tiny_problems.dag.nodes.log",
1162 "reldir/001/tiny_problems.node_status",
1163 },
1164 )
1166 def testSubdags(self):
1167 with temporaryDirectory() as tmp_dir:
1168 test_tmp_dir = pathlib.Path(tmp_dir)
1169 submit_dir = test_tmp_dir / "submit"
1170 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1171 lssthtc.htc_backup_files(submit_dir)
1172 result_submit = []
1173 for root, _, files in os.walk(submit_dir):
1174 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1175 self.assertEqual(
1176 set(result_submit),
1177 {
1178 "./group_failed_1.dag",
1179 "./group_failed_1.dag.dagman.log",
1180 "./group_failed_1.dag.dagman.out",
1181 "./group_failed_1.dag.rescue001",
1182 "subdags/wms_group_order1_val1a/group_order1_val1a.dag",
1183 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.dagman.log",
1184 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.dagman.out",
1185 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.nodes.log",
1186 "subdags/wms_group_order1_val1a/group_order1_val1a.node_status",
1187 "subdags/wms_group_order1_val1a/wms_group_order1_val1a.dag.post.out",
1188 "subdags/wms_group_order1_val1a/wms_group_order1_val1a.status.txt",
1189 "subdags/wms_group_order1_val1b/group_order1_val1b.dag",
1190 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.dagman.log",
1191 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.dagman.out",
1192 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.rescue001",
1193 "subdags/wms_group_order1_val1c/group_order1_val1c.dag",
1194 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.dagman.log",
1195 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.dagman.out",
1196 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.nodes.log",
1197 "subdags/wms_group_order1_val1c/group_order1_val1c.node_status",
1198 "subdags/wms_group_order1_val1c/wms_group_order1_val1c.dag.post.out",
1199 "subdags/wms_group_order1_val1c/wms_group_order1_val1c.status.txt",
1200 "001/group_failed_1.dag.nodes.log",
1201 "001/group_failed_1.info.json",
1202 "001/group_failed_1.node_status",
1203 "001/subdags/wms_group_order1_val1b/group_order1_val1b.dag.nodes.log",
1204 "001/subdags/wms_group_order1_val1b/group_order1_val1b.node_status",
1205 "001/subdags/wms_group_order1_val1b/wms_group_order1_val1b.status.txt",
1206 "001/subdags/wms_group_order1_val1b/wms_group_order1_val1b.dag.post.out",
1207 },
1208 )
1211class UpdateRescueFileTestCase(unittest.TestCase):
1212 """Test _update_rescue_file function."""
1214 def testSuccess(self):
1215 self.maxDiff = None
1216 with temporaryDirectory() as tmp_dir:
1217 test_tmp_dir = pathlib.Path(tmp_dir)
1218 submit_dir = test_tmp_dir / "submit"
1219 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1220 rescue_file = submit_dir / "group_failed_1.dag.rescue001"
1221 failed_subdags = lssthtc._update_rescue_file(rescue_file)
1222 self.assertEqual(set(failed_subdags), {"wms_group_order1_val1b"})
1223 with open(rescue_file) as fh:
1224 lines = fh.readlines()
1225 results = "".join(lines)
1227 truth = """# Rescue DAG file, created after running
1228# the u_testuser_DM-46294_group_fail_20250310T160455Z.dag DAG file
1229# Created 3/10/2025 16:08:56 UTC
1230# Rescue DAG version: 2.0.1 (partial)
1231#
1232# Total number of Nodes: 26
1233# Nodes premarked DONE: 21
1234# Nodes that failed: 2
1235# wms_group_order1_val1b,finalJob,<ENDLIST>
1237DONE pipetaskInit
1238DONE label1_val1c_val2a
1239DONE label1_val1b_val2b
1240DONE label1_val1b_val2a
1241DONE label1_val1c_val2b
1242DONE label1_val1a_val2a
1243DONE label1_val1a_val2b
1244DONE label3_val1c_val2a
1245DONE label3_val1b_val2b
1246DONE label3_val1b_val2a
1247DONE label3_val1c_val2b
1248DONE label3_val1a_val2a
1249DONE label3_val1a_val2b
1250DONE wms_group_order1_val1a
1251DONE label5_val1a_val2a
1252DONE label5_val1a_val2b
1253DONE wms_group_order1_val1c
1254DONE label5_val1c_val2a
1255DONE label5_val1c_val2b
1256DONE wms_check_status_wms_group_order1_val1a
1257DONE wms_check_status_wms_group_order1_val1c
1258"""
1260 self.assertEqual(results, truth)
1263class ReadRescueHeadersTestCase(unittest.TestCase):
1264 """Test _read_rescue_headers function."""
1266 def testTypical(self):
1267 content = "# Header line 1\n# Header line 2\n\nDONE somenode\n"
1268 result = lssthtc._read_rescue_headers(io.StringIO(content))
1269 self.assertEqual(result, ["# Header line 1", "# Header line 2"])
1271 def testEmptyFile(self):
1272 result = lssthtc._read_rescue_headers(io.StringIO(""))
1273 self.assertEqual(result, [])
1275 def testOnlyHeaderLines(self):
1276 content = "# Line 1\n# Line 2\n# Line 3\n"
1277 result = lssthtc._read_rescue_headers(io.StringIO(content))
1278 self.assertEqual(result, ["# Line 1", "# Line 2", "# Line 3"])
1280 def testFirstLineNotComment(self):
1281 content = "DONE somenode\n# Header\n"
1282 result = lssthtc._read_rescue_headers(io.StringIO(content))
1283 self.assertEqual(result, [])
1285 def testWhitespaceStripped(self):
1286 content = " # Header line 1 \n # Header line 2 \n\n"
1287 result = lssthtc._read_rescue_headers(io.StringIO(content))
1288 self.assertEqual(result, ["# Header line 1", "# Header line 2"])
1291class WriteRescueHeadersTestCase(unittest.TestCase):
1292 """Test _write_rescue_headers function."""
1294 def testTypical(self):
1295 header_lines = ["# Header line 1", "# Header line 2", "# Header line 3"]
1296 outfh = io.StringIO()
1297 lssthtc._write_rescue_headers(header_lines, outfh)
1298 self.assertEqual(outfh.getvalue(), "# Header line 1\n# Header line 2\n# Header line 3\n\n")
1300 def testEmptyList(self):
1301 outfh = io.StringIO()
1302 lssthtc._write_rescue_headers([], outfh)
1303 self.assertEqual(outfh.getvalue(), "\n")
1305 def testSingleLine(self):
1306 outfh = io.StringIO()
1307 lssthtc._write_rescue_headers(["# Only line"], outfh)
1308 self.assertEqual(outfh.getvalue(), "# Only line\n\n")
1311class UpdateRescueHeadersTestCase(unittest.TestCase):
1312 """Test _update_rescue_headers function."""
1314 def testWithFailedSubdag(self):
1315 header_lines = [
1316 "# Total number of Nodes: 26",
1317 "# Nodes premarked DONE: 22",
1318 "# Nodes that failed: 2",
1319 "# wms_check_status_wms_group_order1_val1b,finalJob,<ENDLIST>",
1320 ]
1321 result = lssthtc._update_rescue_headers(header_lines)
1322 self.assertEqual(result, ["wms_group_order1_val1b"])
1323 self.assertEqual(header_lines[1], "# Nodes premarked DONE: 21")
1324 self.assertEqual(header_lines[3], "# wms_group_order1_val1b,finalJob,<ENDLIST>")
1326 def testNoSubdagFailures(self):
1327 header_lines = [
1328 "# Nodes premarked DONE: 5",
1329 "# Nodes that failed: 1",
1330 "# finalJob,<ENDLIST>",
1331 ]
1332 result = lssthtc._update_rescue_headers(header_lines)
1333 self.assertEqual(result, [])
1334 self.assertEqual(header_lines[0], "# Nodes premarked DONE: 5")
1335 self.assertEqual(header_lines[2], "# finalJob,<ENDLIST>")
1337 def testMultipleFailedSubdags(self):
1338 header_lines = [
1339 "# Nodes premarked DONE: 10",
1340 "# Nodes that failed: 3",
1341 "# wms_check_status_subdag_a,wms_check_status_subdag_b,finalJob,<ENDLIST>",
1342 ]
1343 result = lssthtc._update_rescue_headers(header_lines)
1344 self.assertEqual(result, ["subdag_a", "subdag_b"])
1345 self.assertEqual(header_lines[0], "# Nodes premarked DONE: 8")
1346 self.assertEqual(header_lines[2], "# subdag_a,subdag_b,finalJob,<ENDLIST>")
1348 def testNoFailedNodesLine(self):
1349 header_lines = [
1350 "# Total number of Nodes: 5",
1351 "# Nodes premarked DONE: 5",
1352 ]
1353 original = list(header_lines)
1354 result = lssthtc._update_rescue_headers(header_lines)
1355 self.assertEqual(result, [])
1356 self.assertEqual(header_lines, original)
1358 def testEmptyHeader(self):
1359 result = lssthtc._update_rescue_headers([])
1360 self.assertEqual(result, [])
1363class ReadDagStatusTestCase(unittest.TestCase):
1364 """Test read_dag_status function and read_single_dag_status."""
1366 def testFileMissing(self):
1367 with temporaryDirectory() as tmp_dir:
1368 with self.assertRaisesRegex(FileNotFoundError, "DAGMan node status not found"):
1369 _ = lssthtc.read_dag_status(tmp_dir)
1371 def testRegular(self):
1372 with temporaryDirectory() as tmp_dir:
1373 submit_dir = os.path.join(tmp_dir, "tiny_problems")
1374 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1375 results = lssthtc.read_dag_status(submit_dir)
1376 truth = {
1377 "JobProcsHeld": 0,
1378 "NodesPost": 0,
1379 "JobProcsIdle": 0,
1380 "NodesTotal": 6,
1381 "NodesFailed": 2,
1382 "NodesDone": 3,
1383 "NodesQueued": 0,
1384 "NodesPre": 0,
1385 "NodesFutile": 1,
1386 "NodesUnready": 0,
1387 }
1388 self.assertEqual(results, results | truth)
1390 def testSubdags(self):
1391 """Making sure it gets data from subdag dirs and doesn't
1392 fail if some subdags haven't started running yet.
1393 """
1394 self.maxDiff = None
1395 with temporaryDirectory() as tmp_dir:
1396 submit_dir = os.path.join(tmp_dir, "submit")
1397 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1398 results = lssthtc.read_dag_status(submit_dir)
1399 truth = {
1400 "JobProcsHeld": 0,
1401 "NodesPost": 0,
1402 "JobProcsIdle": 0,
1403 "NodesTotal": 34,
1404 "NodesFailed": 0,
1405 "NodesDone": 17,
1406 "NodesQueued": 3,
1407 "NodesPre": 0,
1408 "NodesFutile": 0,
1409 "NodesUnready": 14,
1410 }
1411 self.assertEqual(results, results | truth)
1414class ReadDagInfoTestCase(unittest.TestCase):
1415 """Test read_dag_info function."""
1417 def testFileMissing(self):
1418 with temporaryDirectory() as tmp_dir:
1419 with self.assertRaisesRegex(FileNotFoundError, "File with DAGMan job information not found in "):
1420 _ = lssthtc.read_dag_info(tmp_dir)
1422 def testSuccess(self):
1423 with temporaryDirectory() as tmp_dir:
1424 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.info.json", tmp_dir)
1425 filename, results = lssthtc.read_dag_info(tmp_dir)
1426 self.assertIn("info.json", str(filename))
1427 self.assertTrue(pathlib.Path(filename).is_file())
1429 truth = {
1430 "test02": {
1431 "9208.0": {
1432 "ClusterId": 9208,
1433 "GlobalJobId": "test02#9208.0#1739465078",
1434 "bps_wms_service": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorService",
1435 "bps_project": "dev",
1436 "bps_payload": "tiny",
1437 "bps_operator": "testuser",
1438 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorWorkflow",
1439 "bps_provisioning_job": "provisioningJob",
1440 "bps_run_quanta": "label1:1;label2:1",
1441 "bps_campaign": "quick",
1442 "bps_runsite": "testpool",
1443 "bps_job_summary": "pipetaskInit:1;label1:1;label2:1;finalJob:1",
1444 "bps_run": "u_testuser_tiny_20250213T164427Z",
1445 "bps_isjob": "True",
1446 }
1447 }
1448 }
1450 self.assertEqual(results, truth)
1452 def testPermissionError(self):
1453 with temporaryDirectory() as tmp_dir:
1454 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.info.json", tmp_dir)
1455 with unittest.mock.patch("lsst.ctrl.bps.htcondor.lssthtc.open") as mocked_open:
1456 mocked_open.side_effect = PermissionError
1457 with self.assertLogs("lsst.ctrl.bps.htcondor", level="DEBUG") as cm:
1458 _, results = lssthtc.read_dag_info(tmp_dir)
1459 self.assertIn("Retrieving DAGMan job information failed:", cm.output[-1])
1460 self.assertEqual({}, results)
1463class HtcWriteCondorFileTestCase(unittest.TestCase):
1464 """Test htc_write_condor_file function."""
1466 def testSuccess(self):
1467 with temporaryDirectory() as tmp_dir:
1468 job_name = "job1"
1469 filename = pathlib.Path(tmp_dir) / f"label1/{job_name}.sub"
1470 job = {
1471 "executable": "$(CTRL_MPEXEC_DIR)/bin/pipetask",
1472 "arguments": "-a -b 2 -c",
1473 "request_memory": "2000",
1474 "environment": "one=1 two=\"2\" three='spacey 'quoted' value'",
1475 "log": f"{job_name}.log",
1476 }
1477 job_attrs = {
1478 "bps_job_name": job_name,
1479 "bps_job_label": "label1",
1480 "bps_job_quanta": "task1:8;task2:8",
1481 }
1482 expected = [
1483 "executable=$(CTRL_MPEXEC_DIR)/bin/pipetask\n",
1484 'arguments="-a -b 2 -c"\n',
1485 "request_memory=2000\n",
1486 "environment=\"one=1 two=\"2\" three='spacey 'quoted' value'\"\n",
1487 f"output={job_name}.$(Cluster).out\n",
1488 f"error={job_name}.$(Cluster).out\n",
1489 f"log={job_name}.log\n",
1490 f'+bps_job_name = "{job_name}"\n',
1491 '+bps_job_label = "label1"\n',
1492 '+bps_job_quanta = "task1:8;task2:8"\n',
1493 "queue\n",
1494 ]
1496 lssthtc.htc_write_condor_file(filename, job_name, job, job_attrs)
1497 with open(filename, encoding="utf-8") as f:
1498 actual = f.readlines()
1500 self.assertEqual(set(actual), set(expected))
1501 self.assertTrue(filename.exists())
1502 # Try to make Submit object from file to find any syntax issues
1503 _ = lssthtc.htc_create_submit_from_file(filename)
1506class HtcCreateSubmitFromDagTestCase(unittest.TestCase):
1507 """Test htc_create_submit_from_dag function."""
1509 @classmethod
1510 def setUpClass(cls):
1511 cls.bindir = None
1512 # htcondor.Submit.from_dag requires condor_dagman executable in path.
1513 if not which("condor_dagman"): # pragma: no cover
1514 cls.bindir = tempfile.TemporaryDirectory()
1515 fake_dagman_exec = pathlib.Path(cls.bindir.name) / "condor_dagman"
1516 with open(fake_dagman_exec, "w") as fh:
1517 print("#!/bin/bash", file=fh)
1518 print("echo fake_condor_dagman $@", file=fh)
1519 print("exit 0", file=fh)
1520 fake_dagman_exec.chmod(fake_dagman_exec.stat().st_mode | stat.S_IEXEC)
1521 os.environ["PATH"] = f"{os.environ['PATH']}:{cls.bindir.name}"
1523 @classmethod
1524 def tearDownClass(cls):
1525 if cls.bindir: 1525 ↛ 1526line 1525 didn't jump to line 1526 because the condition on line 1525 was never true
1526 cls.bindir.cleanup()
1528 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"})
1529 def testMaxIdleEnvVar(self):
1530 with temporaryDirectory() as tmp_dir:
1531 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1532 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1533 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {})
1534 self.assertIn("-MaxIdle 42", submit["arguments"])
1536 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"})
1537 def testMaxIdleInDAGManConfig(self):
1538 with temporaryDirectory() as tmp_dir:
1539 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1540 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1541 config_filename = pathlib.Path(tmp_dir) / "dagman.conf"
1542 with open(config_filename, "w") as fh:
1543 print("DAGMAN_MAX_JOBS_IDLE = 300", file=fh)
1544 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, config_filename)
1545 self.assertIn("-MaxIdle 300", submit["arguments"])
1547 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"})
1548 def testMaxIdleNotInDAGManConfig(self):
1549 with temporaryDirectory() as tmp_dir:
1550 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1551 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1552 config_filename = pathlib.Path(tmp_dir) / "dagman.conf"
1553 with open(config_filename, "w") as fh:
1554 print("DAGMAN_MAX_JOBS_SUBMITTED = 300", file=fh)
1555 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, config_filename)
1556 self.assertIn("-MaxIdle 42", submit["arguments"])
1558 @unittest.mock.patch.dict(os.environ, {})
1559 def testMaxIdleGiven(self):
1560 with temporaryDirectory() as tmp_dir:
1561 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1562 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1563 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {"MaxIdle": 37})
1564 self.assertIn("-MaxIdle 37", submit["arguments"])
1566 @unittest.mock.patch.dict(os.environ, {})
1567 def testMaxJobsIdleParam(self):
1568 def _fake_params_contains(key):
1569 if key == "DAGMAN_MAX_JOBS_IDLE":
1570 return True
1571 return False # pragma: no cover
1573 def _fake_params_get(key):
1574 if key == "DAGMAN_MAX_JOBS_IDLE":
1575 return 16
1576 return "FAKE_VAL" # pragma: no cover
1578 with temporaryDirectory() as tmp_dir:
1579 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1580 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1581 with unittest.mock.patch("htcondor.param") as mock_param:
1582 mock_param.__contains__.side_effect = _fake_params_contains
1583 mock_param.__getitem__.side_effect = _fake_params_get
1584 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, None)
1585 self.assertIn("-MaxIdle 16", submit["arguments"])
1587 @unittest.mock.patch.dict(os.environ, {})
1588 def testNoMaxJobsIdle(self):
1589 """Note: Since the produced arguments differ depending on
1590 HTCondor version when no MaxIdle passed to from_dag, not
1591 checking arguments string here. Instead just making sure
1592 lssthtc code doesn't pass MaxIdle value to from_dag.
1593 """
1594 with temporaryDirectory() as tmp_dir:
1595 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1596 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1597 with unittest.mock.patch("htcondor.Submit.from_dag") as submit_mock:
1598 with unittest.mock.patch("htcondor.param") as mock_param:
1599 mock_param.__contains__.return_value = False
1600 _ = lssthtc.htc_create_submit_from_dag(str(dag_filename), {})
1601 submit_mock.assert_called_once_with(str(dag_filename), {})
1604class HtcDagTestCase(unittest.TestCase):
1605 """Test for HTCDag class."""
1607 def setUp(self):
1608 job = lssthtc.HTCJob(name="test_job")
1609 job.add_job_cmds(
1610 {
1611 "executable": "/usr/bin/echo",
1612 "arguments": "foo",
1613 "output": "test_job.$(Cluster).out",
1614 "error": "test_job.$(Cluster).out",
1615 "log": "test_job.$(Cluster).log",
1616 }
1617 )
1618 job.subfile = f"{job.name}.sub"
1620 self.dag = lssthtc.HTCDag(name="test_workflow")
1621 self.dag.add_job(job)
1623 self.subfile_expected = [
1624 "executable=/usr/bin/echo\n",
1625 'arguments="foo"\n',
1626 "output=test_job.$(Cluster).out\n",
1627 "error=test_job.$(Cluster).out\n",
1628 "log=test_job.$(Cluster).log\n",
1629 "queue\n",
1630 ]
1632 def tearDown(self):
1633 pass
1635 def testWriteWithDagConfig(self):
1636 with temporaryDirectory() as tmp_dir:
1637 config = BpsConfig(Config(htcondor_config.HTC_DEFAULTS_URI))
1638 job = self.dag.nodes["test_job"]["data"]
1639 wms_config_filename = "dagman.conf"
1640 wms_configurator = dagman_configurator.DagmanConfigurator(config)
1641 wms_configurator.prepare(wms_config_filename, prefix=tmp_dir)
1642 wms_configurator.configure(self.dag)
1643 dagfile_expected = [
1644 f"CONFIG {wms_config_filename}\n",
1645 f'JOB {job.name} "{job.subfile}"\n',
1646 f"NODE_STATUS_FILE {self.dag.name}.node_status\n",
1647 f'SET_JOB_ATTR bps_wms_config_path= "{wms_config_filename}"\n',
1648 ]
1650 self.dag.write(tmp_dir, "", "")
1652 self.assertIn("submit_path", self.dag.graph)
1653 self.assertEqual(self.dag.graph["submit_path"], tmp_dir)
1654 self.assertIn("dag_filename", self.dag.graph)
1655 self.assertEqual(self.dag.graph["dag_filename"], f"{self.dag.graph['name']}.dag")
1656 with open(os.path.join(tmp_dir, self.dag.graph["dag_filename"]), encoding="utf-8") as f:
1657 dagfile_actual = f.readlines()
1658 self.assertEqual(dagfile_actual, dagfile_expected)
1659 with open(os.path.join(tmp_dir, job.subfile), encoding="utf-8") as f:
1660 subfile_actual = f.readlines()
1661 self.assertEqual(subfile_actual, self.subfile_expected)
1663 def testWriteWithoutDagConfig(self):
1664 with temporaryDirectory() as tmp_dir:
1665 job = self.dag.nodes["test_job"]["data"]
1666 dagfile_expected = [
1667 f'JOB {job.name} "{job.subfile}"\n',
1668 f"NODE_STATUS_FILE {self.dag.name}.node_status\n",
1669 ]
1671 self.dag.write(tmp_dir, "", "")
1673 self.assertIn("submit_path", self.dag.graph)
1674 self.assertEqual(self.dag.graph["submit_path"], tmp_dir)
1675 self.assertIn("dag_filename", self.dag.graph)
1676 self.assertEqual(self.dag.graph["dag_filename"], f"{self.dag.graph['name']}.dag")
1677 with open(os.path.join(tmp_dir, self.dag.graph["dag_filename"]), encoding="utf-8") as f:
1678 dagfile_actual = f.readlines()
1679 self.assertEqual(dagfile_actual, dagfile_expected)
1680 with open(os.path.join(tmp_dir, job.subfile), encoding="utf-8") as f:
1681 subfile_actual = f.readlines()
1682 self.assertEqual(subfile_actual, self.subfile_expected)
1684 def testWriteLazySubdag(self):
1685 self.maxDiff = None
1686 dag, truth_files = make_lazy_dag("test1", True)
1687 dag.graph["write_dot"] = True
1688 with temporaryDirectory() as tmp_dir:
1689 dag.write(tmp_dir, "", "")
1690 with open(os.path.join(tmp_dir, dag.graph["dag_filename"]), encoding="utf-8") as f:
1691 dagfile_actual = f.readlines()
1692 self.assertIn("DOT test1.dot\n", dagfile_actual)
1694 all_files = []
1695 for root, _, files in pathlib.Path(tmp_dir).walk():
1696 relroot = pathlib.Path(root).relative_to(tmp_dir)
1697 all_files.extend([str(relroot / f) for f in files])
1698 self.assertEqual(sorted(all_files), sorted(truth_files))
1700 @staticmethod
1701 def _make_simple_job(name):
1702 job = lssthtc.HTCJob(name=name)
1703 job.add_job_cmds({"executable": "/usr/bin/echo", "arguments": name})
1704 job.subfile = f"{name}.sub"
1705 return job
1707 def testWriteEdgesAndSpecialJobs(self):
1708 self.maxDiff = None
1709 dag = lssthtc.HTCDag(name="test_special")
1710 dag.add_attribs({"bps_run": "myrun"})
1711 dag.add_job(self._make_simple_job("jobA"))
1712 dag.add_job(self._make_simple_job("jobB"))
1713 dag.add_job_relationships(["jobA"], ["jobB"])
1714 dag.add_final_job(self._make_simple_job("finalJob"))
1715 dag.add_service_job(self._make_simple_job("svc"))
1717 dagfile_expected = [
1718 'JOB jobA "jobA.sub"\n',
1719 'JOB jobB "jobB.sub"\n',
1720 "PARENT jobA CHILD jobB\n",
1721 f"NODE_STATUS_FILE {dag.name}.node_status\n",
1722 'SET_JOB_ATTR bps_run= "myrun"\n',
1723 'FINAL finalJob "finalJob.sub"\n',
1724 'SERVICE svc "svc.sub"\n',
1725 ]
1727 with temporaryDirectory() as tmp_dir:
1728 dag.write(tmp_dir, "", "")
1729 with open(os.path.join(tmp_dir, dag.graph["dag_filename"]), encoding="utf-8") as f:
1730 self.assertEqual(f.readlines(), dagfile_expected)
1731 # Submit files are written for regular and special jobs.
1732 for subfile in ("jobA.sub", "jobB.sub", "finalJob.sub", "svc.sub"):
1733 self.assertTrue(os.path.exists(os.path.join(tmp_dir, subfile)))
1735 def testWriteDagSubdir(self):
1736 # A non-empty dag_subdir places the .dag file (and DIR clause) under
1737 # that subdirectory, creating it as needed.
1738 dag = lssthtc.HTCDag(name="test_subdir")
1739 dag.add_job(self._make_simple_job("jobA"))
1741 with temporaryDirectory() as tmp_dir:
1742 dag.write(tmp_dir, "", "daglevel")
1743 self.assertEqual(dag.graph["dag_filename"], "daglevel/test_subdir.dag")
1744 self.assertTrue(os.path.exists(os.path.join(tmp_dir, "daglevel", "test_subdir.dag")))
1746 def testWriteMissingJobData(self):
1747 # A node without the "data" key should raise KeyError.
1748 dag = lssthtc.HTCDag(name="test_baddata")
1749 dag.add_node("orphan") # no data=... provided
1750 with temporaryDirectory() as tmp_dir:
1751 with self.assertRaises(KeyError):
1752 dag.write(tmp_dir, "", "")
1755class WriteDagInfoTestCase(unittest.TestCase):
1756 """Test for write_dag_info function."""
1758 def setUp(self):
1759 self.run = "u_testuser_DM-53494_20260220T001651Z"
1760 self.data = {
1761 "mycomputer": {
1762 "24390.0": {
1763 "ClusterId": 24390,
1764 "GlobalJobId": "mycomputer#24390.0#1771546612",
1765 "bps_run": self.run,
1766 "bps_isjob": "True",
1767 "bps_payload": "DM-53494",
1768 "bps_project": "dev",
1769 "bps_runsite": "site1",
1770 "bps_campaign": "ci_rc2",
1771 "bps_operator": "testuser",
1772 "bps_run_quanta": "",
1773 "bps_job_summary": "buildQuantumGraph:1;preparePayloadWorkflow:1;dummyJob:1",
1774 "bps_wms_service": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorService",
1775 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_workflow.HTCondorWorkflow",
1776 "bps_wms_config_path": "dagman.conf",
1777 }
1778 }
1779 }
1781 def testWrite(self):
1782 with temporaryDirectory() as tmp_dir:
1783 with chdir(tmp_dir):
1784 path = pathlib.Path(tmp_dir) / "test.info.json"
1785 filename = lssthtc.write_dag_info(path, self.data)
1786 self.assertTrue(path.is_file(), f"File not found at {path}")
1787 self.assertEqual(filename, path)
1789 read_filename, read_data = lssthtc.read_dag_info(tmp_dir)
1790 self.assertEqual(read_filename, path)
1791 self.assertEqual(read_data, self.data)
1794if __name__ == "__main__":
1795 unittest.main()