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