Coverage for tests/test_lssthtc.py: 99%
761 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-21 05:19 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-21 05:19 +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
40from lsst.ctrl.bps import BpsConfig
41from lsst.ctrl.bps.htcondor import dagman_configurator, htcondor_config, lssthtc
42from lsst.daf.butler import Config
43from lsst.utils.tests import temporaryDirectory
45logger = logging.getLogger("lsst.ctrl.bps.htcondor")
46TESTDIR = os.path.abspath(os.path.dirname(__file__))
49class TestLsstHtc(unittest.TestCase):
50 """Test basic usage."""
52 def testHtcEscapeInt(self):
53 self.assertEqual(lssthtc.htc_escape(100), 100)
55 def testHtcEscapeDouble(self):
56 self.assertEqual(lssthtc.htc_escape('"double"'), '""double""')
58 def testHtcEscapeSingle(self):
59 self.assertEqual(lssthtc.htc_escape("'single'"), "''single''")
61 def testHtcEscapeNoSideEffect(self):
62 val = "'val'"
63 self.assertEqual(lssthtc.htc_escape(val), "''val''")
64 self.assertEqual(val, "'val'")
66 def testHtcEscapeQuot(self):
67 self.assertEqual(lssthtc.htc_escape(""val""), '"val"')
69 def testHtcVersion(self):
70 ver = lssthtc.htc_version()
71 self.assertRegex(ver, r"^\d+\.\d+\.\d+$")
74class HtcTweakJobInfoTestCase(unittest.TestCase):
75 """Test the function responsible for massaging job information."""
77 def setUp(self):
78 self.log_dir = tempfile.TemporaryDirectory()
79 self.log_dirname = pathlib.Path(self.log_dir.name)
80 self.job = {
81 "Cluster": 1,
82 "Proc": 0,
83 "Iwd": str(self.log_dirname),
84 "Owner": self.log_dirname.owner(),
85 "MyType": None,
86 "TerminatedNormally": True,
87 }
89 def tearDown(self):
90 self.log_dir.cleanup()
92 def testDirectAssignments(self):
93 lssthtc.htc_tweak_log_info(self.log_dirname, self.job)
94 self.assertEqual(self.job["ClusterId"], self.job["Cluster"])
95 self.assertEqual(self.job["ProcId"], self.job["Proc"])
96 self.assertEqual(self.job["Iwd"], str(self.log_dirname))
97 self.assertEqual(self.job["Owner"], self.log_dirname.owner())
99 def testIncompatibleAdPassThru(self):
100 # Passing a job ad with insufficient information should be a no-op.
101 expected = {"foo": "bar"}
102 result = dict(expected)
103 lssthtc.htc_tweak_log_info(self.log_dirname, result)
104 self.assertEqual(result, expected)
106 def testJobStatusAssignmentJobAbortedEvent(self):
107 job = self.job | {"MyType": "JobAbortedEvent"}
108 lssthtc.htc_tweak_log_info(self.log_dirname, job)
109 self.assertTrue("JobStatus" in job)
110 self.assertEqual(job["JobStatus"], htcondor.JobStatus.REMOVED)
112 def testJobStatusAssignmentExecuteEvent(self):
113 job = self.job | {"MyType": "ExecuteEvent"}
114 lssthtc.htc_tweak_log_info(self.log_dirname, job)
115 self.assertTrue("JobStatus" in job)
116 self.assertEqual(job["JobStatus"], htcondor.JobStatus.RUNNING)
118 def testJobStatusAssignmentSubmitEvent(self):
119 job = self.job | {"MyType": "SubmitEvent"}
120 lssthtc.htc_tweak_log_info(self.log_dirname, job)
121 self.assertTrue("JobStatus" in job)
122 self.assertEqual(job["JobStatus"], htcondor.JobStatus.IDLE)
124 def testJobStatusAssignmentJobHeldEvent(self):
125 job = self.job | {"MyType": "JobHeldEvent"}
126 lssthtc.htc_tweak_log_info(self.log_dirname, job)
127 self.assertTrue("JobStatus" in job)
128 self.assertEqual(job["JobStatus"], htcondor.JobStatus.HELD)
130 def testJobStatusAssignmentJobTerminatedEvent(self):
131 job = self.job | {"MyType": "JobTerminatedEvent"}
132 lssthtc.htc_tweak_log_info(self.log_dirname, job)
133 self.assertTrue("JobStatus" in job)
134 self.assertEqual(job["JobStatus"], htcondor.JobStatus.COMPLETED)
136 def testJobStatusAssignmentPostScriptTerminatedEvent(self):
137 job = self.job | {"MyType": "PostScriptTerminatedEvent"}
138 lssthtc.htc_tweak_log_info(self.log_dirname, job)
139 self.assertTrue("JobStatus" in job)
140 self.assertEqual(job["JobStatus"], htcondor.JobStatus.COMPLETED)
142 def testJobStatusAssignmentReleaseEventMainDagJob(self):
143 job = self.job | {"MyType": "JobReleaseEvent"}
144 lssthtc.htc_tweak_log_info(self.log_dirname, job)
145 self.assertTrue("JobStatus" in job)
146 self.assertEqual(job["JobStatus"], htcondor.JobStatus.RUNNING)
148 def testJobStatusAssignmentReleaseEventForNodeJob(self):
149 job = self.job | {"MyType": "JobReleaseEvent", "DAGNodeName": "test_payload_job"}
150 lssthtc.htc_tweak_log_info(self.log_dirname, job)
151 self.assertTrue("JobStatus" in job)
152 self.assertEqual(job["JobStatus"], None)
154 def testAddingExitStatusSuccess(self):
155 job = self.job | {
156 "MyType": "JobTerminatedEvent",
157 "ToE": {"ExitBySignal": False, "ExitCode": 1},
158 }
159 lssthtc.htc_tweak_log_info(self.log_dirname, job)
160 self.assertIn("ExitBySignal", job)
161 self.assertIs(job["ExitBySignal"], False)
162 self.assertIn("ExitCode", job)
163 self.assertEqual(job["ExitCode"], 1)
165 def testAddingExitStatusFailure(self):
166 job = self.job | {
167 "MyType": "JobHeldEvent",
168 }
169 with self.assertLogs(logger=logger, level="ERROR") as cm:
170 lssthtc.htc_tweak_log_info(self.log_dirname, job)
171 self.assertIn("Could not determine exit status", cm.output[0])
173 def testLoggingUnknownLogEvent(self):
174 job = self.job | {"MyType": "Foo"}
175 with self.assertLogs(logger=logger, level="DEBUG") as cm:
176 lssthtc.htc_tweak_log_info(self.log_dirname, job)
177 self.assertIn("Unknown log event", cm.output[1])
179 def testMissingKey(self):
180 job = self.job
181 del job["Cluster"]
182 with self.assertRaises(KeyError) as cm:
183 lssthtc.htc_tweak_log_info(self.log_dirname, job)
184 self.assertEqual(str(cm.exception), "'Cluster'")
187class HtcCheckDagmanOutputTestCase(unittest.TestCase):
188 """Test htc_check_dagman_output function."""
190 def test_missing_output_file(self):
191 with temporaryDirectory() as tmp_dir:
192 with self.assertRaises(FileNotFoundError):
193 _ = lssthtc.htc_check_dagman_output(tmp_dir)
195 def test_permissions_output_file(self):
196 with temporaryDirectory() as tmp_dir:
197 copy2(f"{TESTDIR}/data/test_tmpdir_abort.dag.dagman.out", tmp_dir)
198 os.chmod(f"{tmp_dir}/test_tmpdir_abort.dag.dagman.out", 0o200)
199 print(os.stat(f"{tmp_dir}/test_tmpdir_abort.dag.dagman.out"))
200 results = lssthtc.htc_check_dagman_output(tmp_dir)
201 os.chmod(f"{tmp_dir}/test_tmpdir_abort.dag.dagman.out", 0o600)
202 self.assertIn("Could not read dagman output file", results)
204 def test_submit_failure(self):
205 with temporaryDirectory() as tmp_dir:
206 copy2(f"{TESTDIR}/data/bad_submit.dag.dagman.out", tmp_dir)
207 results = lssthtc.htc_check_dagman_output(tmp_dir)
208 self.assertIn("Warn: Job submission issues (last: ", results)
210 def test_tmpdir_abort(self):
211 with temporaryDirectory() as tmp_dir:
212 copy2(f"{TESTDIR}/data/test_tmpdir_abort.dag.dagman.out", tmp_dir)
213 results = lssthtc.htc_check_dagman_output(tmp_dir)
214 self.assertIn("Cannot submit from /tmp", results)
216 def test_no_messages(self):
217 with temporaryDirectory() as tmp_dir:
218 copy2(f"{TESTDIR}/data/test_no_messages.dag.dagman.out", tmp_dir)
219 results = lssthtc.htc_check_dagman_output(tmp_dir)
220 self.assertEqual("", results)
223class SummarizeDagTestCase(unittest.TestCase):
224 """Test summarize_dag function."""
226 def test_no_dag_file(self):
227 with temporaryDirectory() as tmp_dir:
228 summary, job_name_to_pipetask, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
229 self.assertFalse(len(job_name_to_pipetask))
230 self.assertFalse(len(job_name_to_type))
231 self.assertFalse(summary)
233 def test_success(self):
234 with temporaryDirectory() as tmp_dir:
235 copy2(f"{TESTDIR}/data/good.dag", tmp_dir)
236 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
237 self.assertEqual(summary, "pipetaskInit:1;label1:1;label2:1;label3:1;finalJob:1")
238 self.assertEqual(
239 job_name_to_label,
240 {
241 "pipetaskInit": "pipetaskInit",
242 "0682f8f9-12f0-40a5-971e-8b30c7231e5c_label1_val1_val2": "label1",
243 "d0305e2d-f164-4a85-bd24-06afe6c84ed9_label2_val1_val2": "label2",
244 "2806ecc9-1bba-4362-8fff-ab4e6abb9f83_label3_val1_val2": "label3",
245 "finalJob": "finalJob",
246 },
247 )
248 self.assertEqual(
249 job_name_to_type,
250 {
251 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
252 "0682f8f9-12f0-40a5-971e-8b30c7231e5c_label1_val1_val2": lssthtc.WmsNodeType.PAYLOAD,
253 "d0305e2d-f164-4a85-bd24-06afe6c84ed9_label2_val1_val2": lssthtc.WmsNodeType.PAYLOAD,
254 "2806ecc9-1bba-4362-8fff-ab4e6abb9f83_label3_val1_val2": lssthtc.WmsNodeType.PAYLOAD,
255 "finalJob": lssthtc.WmsNodeType.FINAL,
256 },
257 )
259 def test_service(self):
260 with temporaryDirectory() as tmp_dir:
261 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag", tmp_dir)
262 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
263 self.assertEqual(summary, "pipetaskInit:1;label1:2;label2:2;finalJob:1")
264 self.assertEqual(
265 job_name_to_label,
266 {
267 "pipetaskInit": "pipetaskInit",
268 "057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b": "label1",
269 "4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a": "label1",
270 "40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b": "label2",
271 "696ee50d-e711-40d6-9caf-ee29ae4a656d_label2_val1a_val2a": "label2",
272 "finalJob": "finalJob",
273 "provisioningJob": "provisioningJob",
274 },
275 )
276 self.assertEqual(
277 job_name_to_type,
278 {
279 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
280 "057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
281 "4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
282 "40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
283 "696ee50d-e711-40d6-9caf-ee29ae4a656d_label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
284 "finalJob": lssthtc.WmsNodeType.FINAL,
285 "provisioningJob": lssthtc.WmsNodeType.SERVICE,
286 },
287 )
289 def test_noop(self):
290 with temporaryDirectory() as tmp_dir:
291 copy2(f"{TESTDIR}/data/noop_running_1/noop_running_1.dag", tmp_dir)
292 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(tmp_dir)
293 self.assertEqual(
294 set(summary.split(";")),
295 {"pipetaskInit:1", "label1:6", "label2:6", "label3:6", "label4:6", "label5:6", "finalJob:1"},
296 )
297 self.assertEqual(
298 job_name_to_label,
299 {
300 "label1_val1a_val2a": "label1",
301 "label1_val1a_val2b": "label1",
302 "label1_val1b_val2a": "label1",
303 "label1_val1b_val2b": "label1",
304 "label1_val1c_val2a": "label1",
305 "label1_val1c_val2b": "label1",
306 "label2_val1a_val2a": "label2",
307 "label2_val1a_val2b": "label2",
308 "label2_val1b_val2a": "label2",
309 "label2_val1b_val2b": "label2",
310 "label2_val1c_val2a": "label2",
311 "label2_val1c_val2b": "label2",
312 "label3_val1a_val2a": "label3",
313 "label3_val1a_val2b": "label3",
314 "label3_val1b_val2a": "label3",
315 "label3_val1b_val2b": "label3",
316 "label3_val1c_val2a": "label3",
317 "label3_val1c_val2b": "label3",
318 "label4_val1a_val2a": "label4",
319 "label4_val1a_val2b": "label4",
320 "label4_val1b_val2a": "label4",
321 "label4_val1b_val2b": "label4",
322 "label4_val1c_val2a": "label4",
323 "label4_val1c_val2b": "label4",
324 "label5_val1a_val2a": "label5",
325 "label5_val1a_val2b": "label5",
326 "label5_val1b_val2a": "label5",
327 "label5_val1b_val2b": "label5",
328 "label5_val1c_val2a": "label5",
329 "label5_val1c_val2b": "label5",
330 "finalJob": "finalJob",
331 "pipetaskInit": "pipetaskInit",
332 "wms_noop_order1_val1a": "order1",
333 "wms_noop_order1_val1b": "order1",
334 },
335 )
336 self.assertEqual(
337 job_name_to_type,
338 {
339 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
340 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
341 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
342 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
343 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
344 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
345 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
346 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
347 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
348 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
349 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
350 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
351 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
352 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
353 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
354 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
355 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
356 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
357 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
358 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
359 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
360 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
361 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
362 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
363 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
364 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
365 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
366 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
367 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
368 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
369 "finalJob": lssthtc.WmsNodeType.FINAL,
370 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
371 "wms_noop_order1_val1a": lssthtc.WmsNodeType.NOOP,
372 "wms_noop_order1_val1b": lssthtc.WmsNodeType.NOOP,
373 },
374 )
376 def test_subdags(self):
377 with temporaryDirectory() as tmp_dir:
378 submit_dir = os.path.join(tmp_dir, "group_running_1")
379 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
380 summary, job_name_to_label, job_name_to_type = lssthtc.summarize_dag(submit_dir)
381 self.assertEqual(
382 set(summary.split(";")),
383 {"pipetaskInit:1", "label1:6", "label2:6", "label3:6", "label4:6", "label5:6", "finalJob:1"},
384 )
386 self.assertEqual(
387 job_name_to_label,
388 {
389 "pipetaskInit": "pipetaskInit",
390 "label1_val1b_val2a": "label1",
391 "label1_val1c_val2a": "label1",
392 "label1_val1a_val2b": "label1",
393 "label1_val1b_val2b": "label1",
394 "label1_val1c_val2b": "label1",
395 "label1_val1a_val2a": "label1",
396 "label2_val1a_val2b": "label2",
397 "label2_val1a_val2a": "label2",
398 "label2_val1b_val2a": "label2",
399 "label2_val1b_val2b": "label2",
400 "label2_val1c_val2a": "label2",
401 "label2_val1c_val2b": "label2",
402 "label3_val1b_val2a": "label3",
403 "label3_val1c_val2a": "label3",
404 "label3_val1a_val2b": "label3",
405 "label3_val1b_val2b": "label3",
406 "label3_val1c_val2b": "label3",
407 "label3_val1a_val2a": "label3",
408 "label4_val1a_val2b": "label4",
409 "label4_val1a_val2a": "label4",
410 "label4_val1b_val2a": "label4",
411 "label4_val1b_val2b": "label4",
412 "label4_val1c_val2a": "label4",
413 "label4_val1c_val2b": "label4",
414 "label5_val1a_val2b": "label5",
415 "label5_val1a_val2a": "label5",
416 "label5_val1b_val2a": "label5",
417 "label5_val1b_val2b": "label5",
418 "label5_val1c_val2a": "label5",
419 "label5_val1c_val2b": "label5",
420 "finalJob": "finalJob",
421 "provisioningJob": "provisioningJob",
422 "wms_group_order1_val1a": "order1",
423 "wms_group_order1_val1b": "order1",
424 "wms_group_order1_val1c": "order1",
425 "wms_check_status_wms_group_order1_val1a": "order1",
426 "wms_check_status_wms_group_order1_val1b": "order1",
427 "wms_check_status_wms_group_order1_val1c": "order1",
428 },
429 )
431 self.assertEqual(
432 job_name_to_type,
433 {
434 "pipetaskInit": lssthtc.WmsNodeType.PAYLOAD,
435 "label1_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
436 "label1_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
437 "label1_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
438 "label1_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
439 "label1_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
440 "label1_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
441 "label2_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
442 "label2_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
443 "label2_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
444 "label2_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
445 "label2_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
446 "label2_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
447 "label3_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
448 "label3_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
449 "label3_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
450 "label3_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
451 "label3_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
452 "label3_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
453 "label4_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
454 "label4_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
455 "label4_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
456 "label4_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
457 "label4_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
458 "label4_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
459 "label5_val1a_val2b": lssthtc.WmsNodeType.PAYLOAD,
460 "label5_val1a_val2a": lssthtc.WmsNodeType.PAYLOAD,
461 "label5_val1b_val2a": lssthtc.WmsNodeType.PAYLOAD,
462 "label5_val1b_val2b": lssthtc.WmsNodeType.PAYLOAD,
463 "label5_val1c_val2a": lssthtc.WmsNodeType.PAYLOAD,
464 "label5_val1c_val2b": lssthtc.WmsNodeType.PAYLOAD,
465 "finalJob": lssthtc.WmsNodeType.FINAL,
466 "provisioningJob": lssthtc.WmsNodeType.SERVICE,
467 "wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG,
468 "wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG,
469 "wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG,
470 "wms_check_status_wms_group_order1_val1a": lssthtc.WmsNodeType.SUBDAG_CHECK,
471 "wms_check_status_wms_group_order1_val1b": lssthtc.WmsNodeType.SUBDAG_CHECK,
472 "wms_check_status_wms_group_order1_val1c": lssthtc.WmsNodeType.SUBDAG_CHECK,
473 },
474 )
477class ReadDagNodesLogTestCase(unittest.TestCase):
478 """Test read_dag_nodes_log function."""
480 def setUp(self):
481 self.tmpdir = tempfile.mkdtemp()
483 def tearDown(self):
484 rmtree(self.tmpdir, ignore_errors=True)
486 def testFileMissing(self):
487 with self.assertRaisesRegex(FileNotFoundError, "DAGMan node log not found in"):
488 _ = lssthtc.read_dag_nodes_log(self.tmpdir)
490 def testRegular(self):
491 with temporaryDirectory() as tmp_dir:
492 submit_dir = os.path.join(tmp_dir, "tiny_problems")
493 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
494 results = lssthtc.read_dag_nodes_log(submit_dir)
495 self.assertEqual(results["9231.0"]["Cluster"], 9231)
496 self.assertEqual(results["9231.0"]["Proc"], 0)
497 self.assertEqual(results["9231.0"]["ToE"]["ExitCode"], 1)
498 self.assertEqual(len(results), 6)
500 def testSubdags(self):
501 """Making sure it gets data from subdag dirs and doesn't
502 fail if some subdags haven't started running yet.
503 """
504 with temporaryDirectory() as tmp_dir:
505 submit_dir = os.path.join(tmp_dir, "group_running_1")
506 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
507 results = lssthtc.read_dag_nodes_log(submit_dir)
508 # main dag
509 self.assertEqual(results["10094.0"]["Cluster"], 10094)
510 # subdag
511 self.assertEqual(results["10112.0"]["Cluster"], 10112)
512 self.assertEqual(results["10116.0"]["Cluster"], 10116)
515class ReadNodeStatusTestCase(unittest.TestCase):
516 """Test read_node_status function."""
518 def setUp(self):
519 self.tmpdir = tempfile.mkdtemp()
521 def tearDown(self):
522 rmtree(self.tmpdir, ignore_errors=True)
524 def testServiceJobNotSubmitted(self):
525 # tiny_prov_no_submit files have successful workflow
526 # but provisioningJob could not submit.
527 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag.nodes.log", self.tmpdir)
528 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag.dagman.log", self.tmpdir)
529 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.node_status", self.tmpdir)
530 copy2(f"{TESTDIR}/data/tiny_prov_no_submit/tiny_prov_no_submit.dag", self.tmpdir)
532 jobs = lssthtc.read_node_status(self.tmpdir)
533 found = [
534 id_
535 for id_ in jobs
536 if jobs[id_].get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN) == lssthtc.WmsNodeType.SERVICE
537 ]
538 self.assertEqual(len(found), 1)
539 self.assertEqual(jobs[found[0]]["DAGNodeName"], "provisioningJob")
540 self.assertEqual(jobs[found[0]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
542 def testMissingStatusFile(self):
543 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag.nodes.log", self.tmpdir)
544 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag.dagman.log", self.tmpdir)
545 copy2(f"{TESTDIR}/data/tiny_problems/tiny_problems.dag", self.tmpdir)
547 jobs = lssthtc.read_node_status(self.tmpdir)
548 self.assertEqual(len(jobs), 7)
549 self.assertEqual(jobs["9230.0"]["DAGNodeName"], "pipetaskInit")
550 self.assertEqual(jobs["9230.0"]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD)
551 found = [
552 id_
553 for id_ in jobs
554 if jobs[id_].get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN) == lssthtc.WmsNodeType.SERVICE
555 ]
556 self.assertEqual(len(found), 1)
557 self.assertEqual(jobs[found[0]]["DAGNodeName"], "provisioningJob")
559 def testSubdagsRunning(self):
560 with temporaryDirectory() as tmp_dir:
561 test_tmp_dir = pathlib.Path(tmp_dir)
562 submit_dir = test_tmp_dir / "submit"
563 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
564 jobs = lssthtc.read_node_status(submit_dir)
565 self.assertEqual(len(jobs), 39) # includes non-payload jobs
566 # not guaranteed ids are same, so use names instead
567 job_name_to_id = {}
568 for id_, info in jobs.items():
569 job_name_to_id[info.get("DAGNodeName", id_)] = id_
570 job_type_to_names = {}
571 for id_, info in jobs.items():
572 job_type_to_names.setdefault(
573 info.get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN), set()
574 ).add(info.get("DAGNodeName", id_))
576 # check counts
577 self.assertNotIn(lssthtc.WmsNodeType.NOOP, job_type_to_names)
578 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.PAYLOAD]), 31)
579 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.FINAL]), 1)
580 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SERVICE]), 1)
581 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG]), 3)
582 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG_CHECK]), 3)
584 # spot check some statuses
585 self.assertEqual(
586 jobs[job_name_to_id["label3_val1a_val2b"]]["NodeStatus"], lssthtc.NodeStatus.DONE
587 )
588 self.assertEqual(
589 jobs[job_name_to_id["wms_group_order1_val1a"]]["NodeStatus"], lssthtc.NodeStatus.SUBMITTED
590 )
591 self.assertEqual(
592 jobs[job_name_to_id["label5_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY
593 )
594 self.assertEqual(
595 jobs[job_name_to_id["label2_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.DONE
596 )
598 def testSubdagsFailed(self):
599 with temporaryDirectory() as tmp_dir:
600 test_tmp_dir = pathlib.Path(tmp_dir)
601 submit_dir = test_tmp_dir / "submit"
602 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
603 jobs = lssthtc.read_node_status(submit_dir)
604 self.assertEqual(len(jobs), 39)
605 # not guaranteed ids are same, so use names instead
606 job_name_to_id = {}
607 for id_, info in jobs.items():
608 job_name_to_id[info.get("DAGNodeName", id_)] = id_
609 job_type_to_names = {}
610 for id_, info in jobs.items():
611 job_type_to_names.setdefault(
612 info.get("wms_node_type", lssthtc.WmsNodeType.UNKNOWN), set()
613 ).add(info.get("DAGNodeName", id_))
615 # check counts
616 self.assertNotIn(lssthtc.WmsNodeType.NOOP, job_type_to_names)
617 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.PAYLOAD]), 31)
618 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.FINAL]), 1)
619 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SERVICE]), 1)
620 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG]), 3)
621 self.assertEqual(len(job_type_to_names[lssthtc.WmsNodeType.SUBDAG_CHECK]), 3)
623 # spot check some statuses
624 self.assertEqual(
625 jobs[job_name_to_id["label3_val1a_val2b"]]["NodeStatus"], lssthtc.NodeStatus.DONE
626 )
627 self.assertEqual(
628 jobs[job_name_to_id["wms_group_order1_val1a"]]["NodeStatus"], lssthtc.NodeStatus.DONE
629 )
630 self.assertEqual(
631 jobs[job_name_to_id["label5_val1a_val2a"]]["NodeStatus"], lssthtc.NodeStatus.DONE
632 )
634 self.assertEqual(
635 jobs[job_name_to_id["label5_val1b_val2a"]]["NodeStatus"], lssthtc.NodeStatus.FUTILE
636 )
637 self.assertEqual(
638 jobs[job_name_to_id["wms_group_order1_val1b"]]["NodeStatus"], lssthtc.NodeStatus.DONE
639 )
640 self.assertEqual(
641 jobs[job_name_to_id["wms_check_status_wms_group_order1_val1b"]]["NodeStatus"],
642 lssthtc.NodeStatus.ERROR,
643 )
646class ReadSingleNodeStatusTestCase(unittest.TestCase):
647 """Test read_single_node_status function."""
649 def setUp(self):
650 self.tmpdir = tempfile.mkdtemp()
652 def tearDown(self):
653 rmtree(self.tmpdir, ignore_errors=True)
655 def _copyFiles(self, data_subdir, suffixes):
656 """Copy files with given suffixes from tests/data/<data_subdir>/."""
657 for suffix in suffixes:
658 copy2(f"{TESTDIR}/data/{data_subdir}/{data_subdir}{suffix}", self.tmpdir)
660 def _jobNameToId(self, jobs):
661 return {info["DAGNodeName"]: id_ for id_, info in jobs.items()}
663 def testAllDone(self):
664 self._copyFiles(
665 "tiny_success",
666 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
667 )
668 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
669 jobs = lssthtc.read_single_node_status(filename, -1)
671 self.assertEqual(len(jobs), 5)
672 name_to_id = self._jobNameToId(jobs)
674 # All four submitted nodes are marked DONE.
675 for name in [
676 "pipetaskInit",
677 "5bba27bd-8df7-4668-a9c5-e911192c5cdb_label1_val1_val2",
678 "0b225f1f-6edf-4380-b546-76c97947a88f_label2_val1_val2",
679 "finalJob",
680 ]:
681 self.assertIn(name, name_to_id, msg=f"Missing job {name}")
682 self.assertEqual(
683 jobs[name_to_id[name]]["NodeStatus"],
684 lssthtc.NodeStatus.DONE,
685 msg=f"Expected DONE for {name}",
686 )
688 # Service job not tracked by node_status; it came from the event log so
689 # it has a real positive ClusterId but no NodeStatus field.
690 self.assertIn("provisioningJob", name_to_id)
691 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
693 # Spot-check labels and types.
694 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["bps_job_label"], "pipetaskInit")
695 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD)
696 self.assertEqual(jobs[name_to_id["finalJob"]]["wms_node_type"], lssthtc.WmsNodeType.FINAL)
697 self.assertEqual(jobs[name_to_id["provisioningJob"]]["wms_node_type"], lssthtc.WmsNodeType.SERVICE)
699 # DAGManJobID is populated from the dagman log for every job.
700 for job in jobs.values():
701 self.assertIn("DAGManJobID", job)
703 def testMixedStatuses(self):
704 self._copyFiles(
705 "tiny_problems",
706 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
707 )
708 filename = pathlib.Path(self.tmpdir) / "tiny_problems.node_status"
709 jobs = lssthtc.read_single_node_status(filename, -1)
711 self.assertEqual(len(jobs), 7)
712 name_to_id = self._jobNameToId(jobs)
714 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
715 self.assertEqual(
716 jobs[name_to_id["057c8caf-66f6-4612-abf7-cdea5b666b1b_label1_val1a_val2b"]]["NodeStatus"],
717 lssthtc.NodeStatus.ERROR,
718 )
719 self.assertEqual(
720 jobs[name_to_id["4a7f478b-2e9b-435c-a730-afac3f621658_label1_val1a_val2a"]]["NodeStatus"],
721 lssthtc.NodeStatus.DONE,
722 )
723 self.assertEqual(
724 jobs[name_to_id["40040b97-606d-4997-98d3-e0493055fe7e_label2_val1a_val2b"]]["NodeStatus"],
725 lssthtc.NodeStatus.FUTILE,
726 )
727 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.ERROR)
728 # Service job not tracked by node_status; came from event log so
729 # it has a real positive ClusterId but no NodeStatus field.
730 self.assertIn("provisioningJob", name_to_id)
731 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
733 def testRunningWorkflow(self):
734 self._copyFiles(
735 "tiny_running",
736 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
737 )
738 filename = pathlib.Path(self.tmpdir) / "tiny_running.node_status"
739 jobs = lssthtc.read_single_node_status(filename, -1)
741 self.assertEqual(len(jobs), 5)
742 name_to_id = self._jobNameToId(jobs)
744 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
745 self.assertEqual(
746 jobs[name_to_id["ca27ea57-c014-44c1-838a-78c06bc3ec1b_label1_val1_val2"]]["NodeStatus"],
747 lssthtc.NodeStatus.SUBMITTED,
748 )
749 self.assertEqual(
750 jobs[name_to_id["dbf919fa-5453-4b05-8806-ad6390fda0a3_label2_val1_val2"]]["NodeStatus"],
751 lssthtc.NodeStatus.NOT_READY,
752 )
753 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
754 # Service job appeared in the event log; has a real positive ClusterId.
755 self.assertIn("provisioningJob", name_to_id)
756 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
758 def testMissingNodeStatusFile(self):
759 # Omit the .node_status file; jobs must be built from the event log
760 # and dag.
761 self._copyFiles(
762 "tiny_problems",
763 [".dag", ".dag.dagman.log", ".dag.nodes.log"],
764 )
765 filename = pathlib.Path(self.tmpdir) / "tiny_problems.node_status"
766 jobs = lssthtc.read_single_node_status(filename, -1)
768 self.assertEqual(len(jobs), 7)
769 name_to_id = self._jobNameToId(jobs)
771 # Jobs that appeared in the event log have real (positive) cluster IDs.
772 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["DAGNodeName"], "pipetaskInit")
773 self.assertGreater(jobs[name_to_id["pipetaskInit"]]["ClusterId"], 0)
775 # The service job appeared in the event log and has a real positive ID.
776 self.assertGreater(jobs[name_to_id["provisioningJob"]]["ClusterId"], 0)
778 # All jobs carry the correct label and type from the dag file.
779 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["wms_node_type"], lssthtc.WmsNodeType.PAYLOAD)
780 self.assertEqual(jobs[name_to_id["provisioningJob"]]["wms_node_type"], lssthtc.WmsNodeType.SERVICE)
782 def testMissingLogFiles(self):
783 # Omit both log files; every job should get a fake negative ID.
784 self._copyFiles("tiny_success", [".dag", ".node_status"])
785 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
786 jobs = lssthtc.read_single_node_status(filename, -1)
788 self.assertEqual(len(jobs), 5)
789 for job in jobs.values():
790 self.assertLess(job["ClusterId"], 0)
792 # NodeStatus values from the node_status file must still be preserved.
793 name_to_id = self._jobNameToId(jobs)
794 self.assertEqual(jobs[name_to_id["pipetaskInit"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
795 self.assertEqual(jobs[name_to_id["finalJob"]]["NodeStatus"], lssthtc.NodeStatus.DONE)
796 self.assertEqual(jobs[name_to_id["provisioningJob"]]["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
798 def testInitFakeId(self):
799 # Verify fake IDs count down from the given starting value.
800 self._copyFiles("tiny_success", [".dag", ".node_status"])
801 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
802 init_fake_id = -10
803 jobs = lssthtc.read_single_node_status(filename, init_fake_id)
805 self.assertEqual(len(jobs), 5)
806 cluster_ids = [job["ClusterId"] for job in jobs.values()]
807 # All IDs must be at most init_fake_id (i.e., -10 or lower).
808 for cid in cluster_ids:
809 self.assertLessEqual(cid, init_fake_id)
810 # All IDs must be unique.
811 self.assertEqual(len(set(cluster_ids)), len(cluster_ids))
813 def testFromDagJobAttribute(self):
814 self._copyFiles(
815 "tiny_success",
816 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
817 )
818 filename = pathlib.Path(self.tmpdir) / "tiny_success.node_status"
819 jobs = lssthtc.read_single_node_status(filename, -1)
821 for id_, job in jobs.items():
822 self.assertEqual(
823 job["from_dag_job"],
824 "wms_tiny_success",
825 msg=f"Job {id_} has wrong from_dag_job",
826 )
828 def testServiceJobPlaceholder(self):
829 self._copyFiles(
830 "tiny_prov_no_submit",
831 [".dag", ".dag.dagman.log", ".dag.nodes.log", ".node_status"],
832 )
833 filename = pathlib.Path(self.tmpdir) / "tiny_prov_no_submit.node_status"
834 jobs = lssthtc.read_single_node_status(filename, -1)
836 service_jobs = [
837 (id_, info)
838 for id_, info in jobs.items()
839 if info.get("wms_node_type") == lssthtc.WmsNodeType.SERVICE
840 ]
841 self.assertEqual(len(service_jobs), 1)
842 service_id, service_job = service_jobs[0]
843 self.assertEqual(service_job["DAGNodeName"], "provisioningJob")
844 self.assertEqual(service_job["NodeStatus"], lssthtc.NodeStatus.NOT_READY)
845 self.assertLess(service_job["ClusterId"], 0)
848class HTCJobTestCase(unittest.TestCase):
849 """Test HTCJob methods."""
851 def testWriteDagCommandsPayload(self):
852 job = lssthtc.HTCJob(
853 "job1",
854 "label1",
855 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
856 {"dir": "jobs/label1"},
857 )
858 job.subfile = "job1.sub"
860 mockfh = io.StringIO()
861 job.write_dag_commands(mockfh, "../..")
862 self.assertIn('JOB job1 "job1.sub" DIR "../../jobs/label1"', mockfh.getvalue())
864 def testWriteDagCommandsNotJob(self):
865 # Testing giving command_name, no dag_rel_path and no dir
866 job = lssthtc.HTCJob(
867 "finalJob",
868 "finalJob",
869 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
870 )
871 job.subfile = "jobs/finalJob/finalJob.sub"
872 mockfh = io.StringIO()
873 job.write_dag_commands(mockfh, "", "FINAL")
874 self.assertIn('FINAL finalJob "jobs/finalJob/finalJob.sub"', mockfh.getvalue())
876 def testWriteDagCommandsNoop(self):
877 job = lssthtc.HTCJob("wms_noop_job1", "label1", {}, {"noop": True})
878 job.subfile = "notthere.sub"
879 mockfh = io.StringIO()
880 job.write_dag_commands(mockfh, "")
881 self.assertIn("NOOP", mockfh.getvalue())
883 def testWriteSubmitFile(self):
884 job = lssthtc.HTCJob(
885 "job1",
886 "label1",
887 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
888 )
889 with temporaryDirectory() as tmp_dir:
890 filename = pathlib.Path(tmp_dir) / "label1/job1.sub"
891 job.write_submit_file(filename.parent)
892 self.assertTrue(filename.exists())
893 # Try to make Submit object from file to find any syntax issues
894 _ = lssthtc.htc_create_submit_from_file(filename)
896 def testWriteSubmitFileExists(self):
897 job = lssthtc.HTCJob(
898 "job1",
899 "label1",
900 {"executable": "/bin/sleep", "arguments": "60", "log": "job1.log"},
901 )
902 with temporaryDirectory() as tmp_dir:
903 filename = pathlib.Path(tmp_dir) / "job1.sub"
904 job.subfile = filename
905 with open(filename, "w"):
906 pass # make empty file
907 job.write_submit_file(filename.parent)
908 # make sure didn't overwrite file
909 self.assertEqual(filename.stat().st_size, 0, "Incorrectly overwrote existing file")
912class HtcWriteJobCommands(unittest.TestCase):
913 """Test _htc_write_job_commands function."""
915 def testAllCommands(self):
916 dag_cmds = {
917 "pre": {
918 "defer": {"status": 1, "time": 120},
919 "debug": {"filename": "debug_pre.txt", "type": "ALL"},
920 "executable": "exec1",
921 "arguments": "arg1 arg2",
922 },
923 "post": {
924 "defer": {"status": 2, "time": 180},
925 "debug": {"filename": "debug_post.txt", "type": "ALL"},
926 "executable": "exec2",
927 "arguments": "arg3 arg4",
928 },
929 "vars": {"num": 8, "spaces": "a space"},
930 "pre_skip": "1",
931 "retry": 3,
932 "retry_unless_exit": 1,
933 "abort_dag_on": {"node_exit": 100, "abort_exit": 4},
934 "priority": 123,
935 }
937 truth = """SCRIPT DEFER 1 120 DEBUG debug_pre.txt ALL PRE job1 exec1 arg1 arg2
938SCRIPT DEFER 2 180 DEBUG debug_post.txt ALL POST job1 exec2 arg3 arg4
939VARS job1 num="8"
940VARS job1 spaces="a space"
941PRE_SKIP job1 1
942RETRY job1 3 UNLESS-EXIT 1
943ABORT-DAG-ON job1 100 RETURN 4
944PRIORITY job1 123
945"""
946 mockfh = io.StringIO()
947 lssthtc._htc_write_job_commands(mockfh, "job1", dag_cmds)
948 self.assertEqual(mockfh.getvalue(), truth)
950 def testPartialCommands(self):
951 # Trigger skipping the inner if clauses.
952 dag_cmds = {
953 "pre": {
954 "executable": "exec1",
955 },
956 "post": {
957 "executable": "exec2",
958 },
959 "vars": {"num": 8, "spaces": "a space"},
960 "pre_skip": "1",
961 "retry": 3,
962 }
964 truth = """SCRIPT PRE job1 exec1
965SCRIPT POST job1 exec2
966VARS job1 num="8"
967VARS job1 spaces="a space"
968PRE_SKIP job1 1
969RETRY job1 3
970"""
971 mockfh = io.StringIO()
972 lssthtc._htc_write_job_commands(mockfh, "job1", dag_cmds)
973 self.assertEqual(mockfh.getvalue(), truth)
975 def testNoCommands(self):
976 dag_cmds = {}
977 mockfh = io.StringIO()
978 lssthtc._htc_write_job_commands(mockfh, "job2", dag_cmds)
979 self.assertEqual(mockfh.getvalue(), "")
981 def testFinal(self):
982 self.maxDiff = None
983 dag_cmds = {
984 "pre": {
985 "defer": {"status": 1, "time": 120},
986 "debug": {"filename": "debug_pre.txt", "type": "ALL"},
987 "executable": "exec1",
988 "arguments": "arg1 arg2",
989 },
990 "post": {
991 "defer": {"status": 2, "time": 180},
992 "debug": {"filename": "debug_post.txt", "type": "ALL"},
993 "executable": "exec2",
994 "arguments": "arg3 arg4",
995 },
996 "vars": {"num": 8, "spaces": "a space"},
997 "pre_skip": "1",
998 "retry": 3,
999 "retry_unless_exit": 1,
1000 "abort_dag_on": {"node_exit": 100, "abort_exit": 4},
1001 "priority": 123,
1002 }
1004 truth = """SCRIPT DEFER 1 120 DEBUG debug_pre.txt ALL PRE finalJob exec1 arg1 arg2
1005SCRIPT DEFER 2 180 DEBUG debug_post.txt ALL POST finalJob exec2 arg3 arg4
1006VARS finalJob num="8"
1007VARS finalJob spaces="a space"
1008PRE_SKIP finalJob 1
1009"""
1010 mockfh = io.StringIO()
1011 lssthtc._htc_write_job_commands(mockfh, "finalJob", dag_cmds, "FINAL")
1012 self.assertEqual(mockfh.getvalue(), truth)
1015class HTCBackupFilesSinglePathTestCase(unittest.TestCase):
1016 """Test htc_backup_files_single_path function."""
1018 def testSrcDestSame(self):
1019 with temporaryDirectory() as tmp_dir:
1020 with self.assertRaisesRegex(
1021 RuntimeError, "Destination directory is same as the source directory"
1022 ):
1023 lssthtc.htc_backup_files_single_path(tmp_dir, tmp_dir)
1025 def testSuccess(self):
1026 with temporaryDirectory() as tmp_dir:
1027 test_tmp_dir = pathlib.Path(tmp_dir)
1028 submit_dir = test_tmp_dir / "the_src_dir"
1029 copytree(f"{TESTDIR}/data/tiny_success", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1030 backup_dir = test_tmp_dir / "the_dest_dir"
1031 backup_dir.mkdir()
1032 lssthtc.htc_backup_files_single_path(submit_dir, backup_dir)
1033 result_submit = []
1034 for root, _, files in os.walk(submit_dir):
1035 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1036 self.assertEqual(
1037 set(result_submit),
1038 {
1039 "./tiny_success.dag.dagman.log",
1040 "./tiny_success.dag.dagman.out",
1041 "./tiny_success.dag",
1042 },
1043 )
1044 result_backup = []
1045 for root, _, files in os.walk(backup_dir):
1046 result_backup.extend([str(os.path.join(os.path.relpath(root, backup_dir), f)) for f in files])
1047 self.assertEqual(
1048 set(result_backup),
1049 {
1050 "./tiny_success.info.json",
1051 "./tiny_success.dag.metrics",
1052 "./tiny_success.dag.nodes.log",
1053 "./tiny_success.node_status",
1054 },
1055 )
1058class HTCBackupFilesTestCase(unittest.TestCase):
1059 """Test htc_backup_files function."""
1061 def testDirectoryNotFound(self):
1062 with temporaryDirectory() as tmp_dir:
1063 test_tmp_dir = pathlib.Path(tmp_dir)
1064 submit_dir = test_tmp_dir / "submit"
1065 with self.assertRaises(FileNotFoundError):
1066 lssthtc.htc_backup_files(submit_dir)
1068 def testSuccess(self):
1069 with temporaryDirectory() as tmp_dir:
1070 test_tmp_dir = pathlib.Path(tmp_dir)
1071 submit_dir = test_tmp_dir / "submit"
1072 copytree(f"{TESTDIR}/data/tiny_success", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1073 lssthtc.htc_backup_files(submit_dir)
1074 result_submit = []
1075 for root, _, files in os.walk(submit_dir):
1076 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1077 self.assertEqual(
1078 set(result_submit),
1079 {
1080 "./tiny_success.dag.dagman.log",
1081 "./tiny_success.dag.dagman.out",
1082 "./tiny_success.dag",
1083 "000/tiny_success.info.json",
1084 "000/tiny_success.dag.metrics",
1085 "000/tiny_success.dag.nodes.log",
1086 "000/tiny_success.node_status",
1087 },
1088 )
1090 def testDestNotInSubmitDir(self):
1091 with temporaryDirectory() as tmp_dir:
1092 test_tmp_dir = pathlib.Path(tmp_dir)
1093 submit_dir = test_tmp_dir / "submit"
1094 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1095 with self.assertLogs("lsst.ctrl.bps.htcondor", level="WARNING") as cm:
1096 lssthtc.htc_backup_files(submit_dir, test_tmp_dir / "backup")
1097 self.assertIn("Invalid backup location:", cm.output[-1])
1098 lssthtc.htc_backup_files(submit_dir)
1099 result_submit = []
1100 for root, _, files in os.walk(submit_dir):
1101 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1102 self.assertEqual(
1103 set(result_submit),
1104 {
1105 "./tiny_problems.dag.dagman.log",
1106 "./tiny_problems.dag.dagman.out",
1107 "./tiny_problems.dag",
1108 "./tiny_problems.dag.rescue001",
1109 "001/tiny_problems.info.json",
1110 "001/tiny_problems.dag.metrics",
1111 "001/tiny_problems.dag.nodes.log",
1112 "001/tiny_problems.node_status",
1113 },
1114 )
1116 def testDestInSubmitDir(self):
1117 with temporaryDirectory() as tmp_dir:
1118 test_tmp_dir = pathlib.Path(tmp_dir)
1119 submit_dir = test_tmp_dir / "submit"
1120 backup_dir = submit_dir / "subdir"
1121 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1122 lssthtc.htc_backup_files(submit_dir, backup_dir)
1123 result_submit = []
1124 for root, _, files in os.walk(submit_dir):
1125 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1126 self.assertEqual(
1127 set(result_submit),
1128 {
1129 "./tiny_problems.dag.dagman.log",
1130 "./tiny_problems.dag.dagman.out",
1131 "./tiny_problems.dag",
1132 "./tiny_problems.dag.rescue001",
1133 "subdir/001/tiny_problems.info.json",
1134 "subdir/001/tiny_problems.dag.metrics",
1135 "subdir/001/tiny_problems.dag.nodes.log",
1136 "subdir/001/tiny_problems.node_status",
1137 },
1138 )
1140 def testRelativeSubdir(self):
1141 with temporaryDirectory() as tmp_dir:
1142 test_tmp_dir = pathlib.Path(tmp_dir)
1143 submit_dir = test_tmp_dir / "submit"
1144 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1145 lssthtc.htc_backup_files(submit_dir, "reldir")
1146 result_submit = []
1147 for root, _, files in os.walk(submit_dir):
1148 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1149 self.assertEqual(
1150 set(result_submit),
1151 {
1152 "./tiny_problems.dag.dagman.log",
1153 "./tiny_problems.dag.dagman.out",
1154 "./tiny_problems.dag",
1155 "./tiny_problems.dag.rescue001",
1156 "reldir/001/tiny_problems.info.json",
1157 "reldir/001/tiny_problems.dag.metrics",
1158 "reldir/001/tiny_problems.dag.nodes.log",
1159 "reldir/001/tiny_problems.node_status",
1160 },
1161 )
1163 def testSubdags(self):
1164 with temporaryDirectory() as tmp_dir:
1165 test_tmp_dir = pathlib.Path(tmp_dir)
1166 submit_dir = test_tmp_dir / "submit"
1167 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1168 lssthtc.htc_backup_files(submit_dir)
1169 result_submit = []
1170 for root, _, files in os.walk(submit_dir):
1171 result_submit.extend([str(os.path.join(os.path.relpath(root, submit_dir), f)) for f in files])
1172 self.assertEqual(
1173 set(result_submit),
1174 {
1175 "./group_failed_1.dag",
1176 "./group_failed_1.dag.dagman.log",
1177 "./group_failed_1.dag.dagman.out",
1178 "./group_failed_1.dag.rescue001",
1179 "subdags/wms_group_order1_val1a/group_order1_val1a.dag",
1180 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.dagman.log",
1181 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.dagman.out",
1182 "subdags/wms_group_order1_val1a/group_order1_val1a.dag.nodes.log",
1183 "subdags/wms_group_order1_val1a/group_order1_val1a.node_status",
1184 "subdags/wms_group_order1_val1a/wms_group_order1_val1a.dag.post.out",
1185 "subdags/wms_group_order1_val1a/wms_group_order1_val1a.status.txt",
1186 "subdags/wms_group_order1_val1b/group_order1_val1b.dag",
1187 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.dagman.log",
1188 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.dagman.out",
1189 "subdags/wms_group_order1_val1b/group_order1_val1b.dag.rescue001",
1190 "subdags/wms_group_order1_val1c/group_order1_val1c.dag",
1191 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.dagman.log",
1192 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.dagman.out",
1193 "subdags/wms_group_order1_val1c/group_order1_val1c.dag.nodes.log",
1194 "subdags/wms_group_order1_val1c/group_order1_val1c.node_status",
1195 "subdags/wms_group_order1_val1c/wms_group_order1_val1c.dag.post.out",
1196 "subdags/wms_group_order1_val1c/wms_group_order1_val1c.status.txt",
1197 "001/group_failed_1.dag.nodes.log",
1198 "001/group_failed_1.info.json",
1199 "001/group_failed_1.node_status",
1200 "001/subdags/wms_group_order1_val1b/group_order1_val1b.dag.nodes.log",
1201 "001/subdags/wms_group_order1_val1b/group_order1_val1b.node_status",
1202 "001/subdags/wms_group_order1_val1b/wms_group_order1_val1b.status.txt",
1203 "001/subdags/wms_group_order1_val1b/wms_group_order1_val1b.dag.post.out",
1204 },
1205 )
1208class UpdateRescueFileTestCase(unittest.TestCase):
1209 """Test _update_rescue_file function."""
1211 def testSuccess(self):
1212 self.maxDiff = None
1213 with temporaryDirectory() as tmp_dir:
1214 test_tmp_dir = pathlib.Path(tmp_dir)
1215 submit_dir = test_tmp_dir / "submit"
1216 copytree(f"{TESTDIR}/data/group_failed_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1217 rescue_file = submit_dir / "group_failed_1.dag.rescue001"
1218 failed_subdags = lssthtc._update_rescue_file(rescue_file)
1219 self.assertEqual(set(failed_subdags), {"wms_group_order1_val1b"})
1220 with open(rescue_file) as fh:
1221 lines = fh.readlines()
1222 results = "".join(lines)
1224 truth = """# Rescue DAG file, created after running
1225# the u_testuser_DM-46294_group_fail_20250310T160455Z.dag DAG file
1226# Created 3/10/2025 16:08:56 UTC
1227# Rescue DAG version: 2.0.1 (partial)
1228#
1229# Total number of Nodes: 26
1230# Nodes premarked DONE: 21
1231# Nodes that failed: 2
1232# wms_group_order1_val1b,finalJob,<ENDLIST>
1234DONE pipetaskInit
1235DONE label1_val1c_val2a
1236DONE label1_val1b_val2b
1237DONE label1_val1b_val2a
1238DONE label1_val1c_val2b
1239DONE label1_val1a_val2a
1240DONE label1_val1a_val2b
1241DONE label3_val1c_val2a
1242DONE label3_val1b_val2b
1243DONE label3_val1b_val2a
1244DONE label3_val1c_val2b
1245DONE label3_val1a_val2a
1246DONE label3_val1a_val2b
1247DONE wms_group_order1_val1a
1248DONE label5_val1a_val2a
1249DONE label5_val1a_val2b
1250DONE wms_group_order1_val1c
1251DONE label5_val1c_val2a
1252DONE label5_val1c_val2b
1253DONE wms_check_status_wms_group_order1_val1a
1254DONE wms_check_status_wms_group_order1_val1c
1255"""
1257 self.assertEqual(results, truth)
1260class ReadRescueHeadersTestCase(unittest.TestCase):
1261 """Test _read_rescue_headers function."""
1263 def testTypical(self):
1264 content = "# Header line 1\n# Header line 2\n\nDONE somenode\n"
1265 result = lssthtc._read_rescue_headers(io.StringIO(content))
1266 self.assertEqual(result, ["# Header line 1", "# Header line 2"])
1268 def testEmptyFile(self):
1269 result = lssthtc._read_rescue_headers(io.StringIO(""))
1270 self.assertEqual(result, [])
1272 def testOnlyHeaderLines(self):
1273 content = "# Line 1\n# Line 2\n# Line 3\n"
1274 result = lssthtc._read_rescue_headers(io.StringIO(content))
1275 self.assertEqual(result, ["# Line 1", "# Line 2", "# Line 3"])
1277 def testFirstLineNotComment(self):
1278 content = "DONE somenode\n# Header\n"
1279 result = lssthtc._read_rescue_headers(io.StringIO(content))
1280 self.assertEqual(result, [])
1282 def testWhitespaceStripped(self):
1283 content = " # Header line 1 \n # Header line 2 \n\n"
1284 result = lssthtc._read_rescue_headers(io.StringIO(content))
1285 self.assertEqual(result, ["# Header line 1", "# Header line 2"])
1288class WriteRescueHeadersTestCase(unittest.TestCase):
1289 """Test _write_rescue_headers function."""
1291 def testTypical(self):
1292 header_lines = ["# Header line 1", "# Header line 2", "# Header line 3"]
1293 outfh = io.StringIO()
1294 lssthtc._write_rescue_headers(header_lines, outfh)
1295 self.assertEqual(outfh.getvalue(), "# Header line 1\n# Header line 2\n# Header line 3\n\n")
1297 def testEmptyList(self):
1298 outfh = io.StringIO()
1299 lssthtc._write_rescue_headers([], outfh)
1300 self.assertEqual(outfh.getvalue(), "\n")
1302 def testSingleLine(self):
1303 outfh = io.StringIO()
1304 lssthtc._write_rescue_headers(["# Only line"], outfh)
1305 self.assertEqual(outfh.getvalue(), "# Only line\n\n")
1308class UpdateRescueHeadersTestCase(unittest.TestCase):
1309 """Test _update_rescue_headers function."""
1311 def testWithFailedSubdag(self):
1312 header_lines = [
1313 "# Total number of Nodes: 26",
1314 "# Nodes premarked DONE: 22",
1315 "# Nodes that failed: 2",
1316 "# wms_check_status_wms_group_order1_val1b,finalJob,<ENDLIST>",
1317 ]
1318 result = lssthtc._update_rescue_headers(header_lines)
1319 self.assertEqual(result, ["wms_group_order1_val1b"])
1320 self.assertEqual(header_lines[1], "# Nodes premarked DONE: 21")
1321 self.assertEqual(header_lines[3], "# wms_group_order1_val1b,finalJob,<ENDLIST>")
1323 def testNoSubdagFailures(self):
1324 header_lines = [
1325 "# Nodes premarked DONE: 5",
1326 "# Nodes that failed: 1",
1327 "# finalJob,<ENDLIST>",
1328 ]
1329 result = lssthtc._update_rescue_headers(header_lines)
1330 self.assertEqual(result, [])
1331 self.assertEqual(header_lines[0], "# Nodes premarked DONE: 5")
1332 self.assertEqual(header_lines[2], "# finalJob,<ENDLIST>")
1334 def testMultipleFailedSubdags(self):
1335 header_lines = [
1336 "# Nodes premarked DONE: 10",
1337 "# Nodes that failed: 3",
1338 "# wms_check_status_subdag_a,wms_check_status_subdag_b,finalJob,<ENDLIST>",
1339 ]
1340 result = lssthtc._update_rescue_headers(header_lines)
1341 self.assertEqual(result, ["subdag_a", "subdag_b"])
1342 self.assertEqual(header_lines[0], "# Nodes premarked DONE: 8")
1343 self.assertEqual(header_lines[2], "# subdag_a,subdag_b,finalJob,<ENDLIST>")
1345 def testNoFailedNodesLine(self):
1346 header_lines = [
1347 "# Total number of Nodes: 5",
1348 "# Nodes premarked DONE: 5",
1349 ]
1350 original = list(header_lines)
1351 result = lssthtc._update_rescue_headers(header_lines)
1352 self.assertEqual(result, [])
1353 self.assertEqual(header_lines, original)
1355 def testEmptyHeader(self):
1356 result = lssthtc._update_rescue_headers([])
1357 self.assertEqual(result, [])
1360class ReadDagStatusTestCase(unittest.TestCase):
1361 """Test read_dag_status function and read_single_dag_status."""
1363 def testFileMissing(self):
1364 with temporaryDirectory() as tmp_dir:
1365 with self.assertRaisesRegex(FileNotFoundError, "DAGMan node status not found"):
1366 _ = lssthtc.read_dag_status(tmp_dir)
1368 def testRegular(self):
1369 with temporaryDirectory() as tmp_dir:
1370 submit_dir = os.path.join(tmp_dir, "tiny_problems")
1371 copytree(f"{TESTDIR}/data/tiny_problems", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1372 results = lssthtc.read_dag_status(submit_dir)
1373 truth = {
1374 "JobProcsHeld": 0,
1375 "NodesPost": 0,
1376 "JobProcsIdle": 0,
1377 "NodesTotal": 6,
1378 "NodesFailed": 2,
1379 "NodesDone": 3,
1380 "NodesQueued": 0,
1381 "NodesPre": 0,
1382 "NodesFutile": 1,
1383 "NodesUnready": 0,
1384 }
1385 self.assertEqual(results, results | truth)
1387 def testSubdags(self):
1388 """Making sure it gets data from subdag dirs and doesn't
1389 fail if some subdags haven't started running yet.
1390 """
1391 self.maxDiff = None
1392 with temporaryDirectory() as tmp_dir:
1393 submit_dir = os.path.join(tmp_dir, "submit")
1394 copytree(f"{TESTDIR}/data/group_running_1", submit_dir, ignore=ignore_patterns("*~", ".???*"))
1395 results = lssthtc.read_dag_status(submit_dir)
1396 truth = {
1397 "JobProcsHeld": 0,
1398 "NodesPost": 0,
1399 "JobProcsIdle": 0,
1400 "NodesTotal": 34,
1401 "NodesFailed": 0,
1402 "NodesDone": 17,
1403 "NodesQueued": 3,
1404 "NodesPre": 0,
1405 "NodesFutile": 0,
1406 "NodesUnready": 14,
1407 }
1408 self.assertEqual(results, results | truth)
1411class ReadDagInfoTestCase(unittest.TestCase):
1412 """Test read_dag_info function."""
1414 def testFileMissing(self):
1415 with temporaryDirectory() as tmp_dir:
1416 with self.assertRaisesRegex(FileNotFoundError, "File with DAGMan job information not found in "):
1417 _ = lssthtc.read_dag_info(tmp_dir)
1419 def testSuccess(self):
1420 with temporaryDirectory() as tmp_dir:
1421 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.info.json", tmp_dir)
1422 results = lssthtc.read_dag_info(tmp_dir)
1424 truth = {
1425 "test02": {
1426 "9208.0": {
1427 "ClusterId": 9208,
1428 "GlobalJobId": "test02#9208.0#1739465078",
1429 "bps_wms_service": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorService",
1430 "bps_project": "dev",
1431 "bps_payload": "tiny",
1432 "bps_operator": "testuser",
1433 "bps_wms_workflow": "lsst.ctrl.bps.htcondor.htcondor_service.HTCondorWorkflow",
1434 "bps_provisioning_job": "provisioningJob",
1435 "bps_run_quanta": "label1:1;label2:1",
1436 "bps_campaign": "quick",
1437 "bps_runsite": "testpool",
1438 "bps_job_summary": "pipetaskInit:1;label1:1;label2:1;finalJob:1",
1439 "bps_run": "u_testuser_tiny_20250213T164427Z",
1440 "bps_isjob": "True",
1441 }
1442 }
1443 }
1445 self.assertEqual(results, truth)
1447 def testPermissionError(self):
1448 with temporaryDirectory() as tmp_dir:
1449 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.info.json", tmp_dir)
1450 with unittest.mock.patch("lsst.ctrl.bps.htcondor.lssthtc.open") as mocked_open:
1451 mocked_open.side_effect = PermissionError
1452 with self.assertLogs("lsst.ctrl.bps.htcondor", level="DEBUG") as cm:
1453 results = lssthtc.read_dag_info(tmp_dir)
1454 self.assertIn("Retrieving DAGMan job information failed:", cm.output[-1])
1455 self.assertEqual({}, results)
1458class HtcWriteCondorFileTestCase(unittest.TestCase):
1459 """Test htc_write_condor_file function."""
1461 def testSuccess(self):
1462 with temporaryDirectory() as tmp_dir:
1463 job_name = "job1"
1464 filename = pathlib.Path(tmp_dir) / f"label1/{job_name}.sub"
1465 job = {
1466 "executable": "$(CTRL_MPEXEC_DIR)/bin/pipetask",
1467 "arguments": "-a -b 2 -c",
1468 "request_memory": "2000",
1469 "environment": "one=1 two=\"2\" three='spacey 'quoted' value'",
1470 "log": f"{job_name}.log",
1471 }
1472 job_attrs = {
1473 "bps_job_name": job_name,
1474 "bps_job_label": "label1",
1475 "bps_job_quanta": "task1:8;task2:8",
1476 }
1477 expected = [
1478 "executable=$(CTRL_MPEXEC_DIR)/bin/pipetask\n",
1479 "arguments=-a -b 2 -c\n",
1480 "request_memory=2000\n",
1481 "environment=\"one=1 two=\"2\" three='spacey 'quoted' value'\"\n",
1482 f"output={job_name}.$(Cluster).out\n",
1483 f"error={job_name}.$(Cluster).out\n",
1484 f"log={job_name}.log\n",
1485 f'+bps_job_name = "{job_name}"\n',
1486 '+bps_job_label = "label1"\n',
1487 '+bps_job_quanta = "task1:8;task2:8"\n',
1488 "queue\n",
1489 ]
1491 lssthtc.htc_write_condor_file(filename, job_name, job, job_attrs)
1492 with open(filename, encoding="utf-8") as f:
1493 actual = f.readlines()
1495 self.assertEqual(set(actual), set(expected))
1496 self.assertTrue(filename.exists())
1497 # Try to make Submit object from file to find any syntax issues
1498 _ = lssthtc.htc_create_submit_from_file(filename)
1501class HtcCreateSubmitFromDagTestCase(unittest.TestCase):
1502 """Test htc_create_submit_from_dag function."""
1504 @classmethod
1505 def setUpClass(cls):
1506 cls.bindir = None
1507 # htcondor.Submit.from_dag requires condor_dagman executable in path.
1508 if not which("condor_dagman"): # pragma: no cover
1509 cls.bindir = tempfile.TemporaryDirectory()
1510 fake_dagman_exec = pathlib.Path(cls.bindir.name) / "condor_dagman"
1511 with open(fake_dagman_exec, "w") as fh:
1512 print("#!/bin/bash", file=fh)
1513 print("echo fake_condor_dagman $@", file=fh)
1514 print("exit 0", file=fh)
1515 fake_dagman_exec.chmod(fake_dagman_exec.stat().st_mode | stat.S_IEXEC)
1516 os.environ["PATH"] = f"{os.environ['PATH']}:{cls.bindir.name}"
1518 @classmethod
1519 def tearDownClass(cls):
1520 if cls.bindir: 1520 ↛ 1521line 1520 didn't jump to line 1521 because the condition on line 1520 was never true
1521 cls.bindir.cleanup()
1523 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"})
1524 def testMaxIdleEnvVar(self):
1525 with temporaryDirectory() as tmp_dir:
1526 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1527 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1528 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {})
1529 self.assertIn("-MaxIdle 42", submit["arguments"])
1531 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"})
1532 def testMaxIdleInDAGManConfig(self):
1533 with temporaryDirectory() as tmp_dir:
1534 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1535 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1536 config_filename = pathlib.Path(tmp_dir) / "dagman.conf"
1537 with open(config_filename, "w") as fh:
1538 print("DAGMAN_MAX_JOBS_IDLE = 300", file=fh)
1539 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, config_filename)
1540 self.assertIn("-MaxIdle 300", submit["arguments"])
1542 @unittest.mock.patch.dict(os.environ, {"_CONDOR_DAGMAN_MAX_JOBS_IDLE": "42"})
1543 def testMaxIdleNotInDAGManConfig(self):
1544 with temporaryDirectory() as tmp_dir:
1545 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1546 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1547 config_filename = pathlib.Path(tmp_dir) / "dagman.conf"
1548 with open(config_filename, "w") as fh:
1549 print("DAGMAN_MAX_JOBS_SUBMITTED = 300", file=fh)
1550 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, config_filename)
1551 self.assertIn("-MaxIdle 42", submit["arguments"])
1553 @unittest.mock.patch.dict(os.environ, {})
1554 def testMaxIdleGiven(self):
1555 with temporaryDirectory() as tmp_dir:
1556 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1557 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1558 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {"MaxIdle": 37})
1559 self.assertIn("-MaxIdle 37", submit["arguments"])
1561 @unittest.mock.patch.dict(os.environ, {})
1562 def testMaxJobsIdleParam(self):
1563 def _fake_params_contains(key):
1564 if key == "DAGMAN_MAX_JOBS_IDLE":
1565 return True
1566 return False # pragma: no cover
1568 def _fake_params_get(key):
1569 if key == "DAGMAN_MAX_JOBS_IDLE":
1570 return 16
1571 return "FAKE_VAL" # pragma: no cover
1573 with temporaryDirectory() as tmp_dir:
1574 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1575 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1576 with unittest.mock.patch("htcondor.param") as mock_param:
1577 mock_param.__contains__.side_effect = _fake_params_contains
1578 mock_param.__getitem__.side_effect = _fake_params_get
1579 submit = lssthtc.htc_create_submit_from_dag(str(dag_filename), {}, None)
1580 self.assertIn("-MaxIdle 16", submit["arguments"])
1582 @unittest.mock.patch.dict(os.environ, {})
1583 def testNoMaxJobsIdle(self):
1584 """Note: Since the produced arguments differ depending on
1585 HTCondor version when no MaxIdle passed to from_dag, not
1586 checking arguments string here. Instead just making sure
1587 lssthtc code doesn't pass MaxIdle value to from_dag.
1588 """
1589 with temporaryDirectory() as tmp_dir:
1590 copy2(f"{TESTDIR}/data/tiny_success/tiny_success.dag", tmp_dir)
1591 dag_filename = pathlib.Path(tmp_dir) / "tiny_success.dag"
1592 with unittest.mock.patch("htcondor.Submit.from_dag") as submit_mock:
1593 with unittest.mock.patch("htcondor.param") as mock_param:
1594 mock_param.__contains__.return_value = False
1595 _ = lssthtc.htc_create_submit_from_dag(str(dag_filename), {})
1596 submit_mock.assert_called_once_with(str(dag_filename), {})
1599class HtcDagTestCase(unittest.TestCase):
1600 """Test for HTCDag class."""
1602 def setUp(self):
1603 job = lssthtc.HTCJob(name="test_job")
1604 job.add_job_cmds(
1605 {
1606 "executable": "/usr/bin/echo",
1607 "arguments": "foo",
1608 "output": "test_job.$(Cluster).out",
1609 "error": "test_job.$(Cluster).out",
1610 "log": "test_job.$(Cluster).log",
1611 }
1612 )
1613 job.subfile = f"{job.name}.sub"
1615 self.dag = lssthtc.HTCDag(name="test_workflow")
1616 self.dag.add_job(job)
1618 self.subfile_expected = [
1619 "executable=/usr/bin/echo\n",
1620 "arguments=foo\n",
1621 "output=test_job.$(Cluster).out\n",
1622 "error=test_job.$(Cluster).out\n",
1623 "log=test_job.$(Cluster).log\n",
1624 "queue\n",
1625 ]
1627 def tearDown(self):
1628 pass
1630 def testWriteWithDagConfig(self):
1631 with temporaryDirectory() as tmp_dir:
1632 config = BpsConfig(Config(htcondor_config.HTC_DEFAULTS_URI))
1633 job = self.dag.nodes["test_job"]["data"]
1634 wms_config_filename = "dagman.conf"
1635 wms_configurator = dagman_configurator.DagmanConfigurator(config)
1636 wms_configurator.prepare(wms_config_filename, prefix=tmp_dir)
1637 wms_configurator.configure(self.dag)
1638 dagfile_expected = [
1639 f"CONFIG {wms_config_filename}\n",
1640 f'JOB {job.name} "{job.subfile}"\n',
1641 f"DOT {self.dag.name}.dot\n",
1642 f"NODE_STATUS_FILE {self.dag.name}.node_status\n",
1643 f'SET_JOB_ATTR bps_wms_config_path= "{wms_config_filename}"\n',
1644 ]
1646 self.dag.write(tmp_dir, "", "")
1648 self.assertIn("submit_path", self.dag.graph)
1649 self.assertEqual(self.dag.graph["submit_path"], tmp_dir)
1650 self.assertIn("dag_filename", self.dag.graph)
1651 self.assertEqual(self.dag.graph["dag_filename"], f"{self.dag.graph['name']}.dag")
1652 with open(os.path.join(tmp_dir, self.dag.graph["dag_filename"]), encoding="utf-8") as f:
1653 dagfile_actual = f.readlines()
1654 self.assertEqual(dagfile_actual, dagfile_expected)
1655 with open(os.path.join(tmp_dir, job.subfile), encoding="utf-8") as f:
1656 subfile_actual = f.readlines()
1657 self.assertEqual(subfile_actual, self.subfile_expected)
1659 def testWriteWithoutDagConfig(self):
1660 with temporaryDirectory() as tmp_dir:
1661 job = self.dag.nodes["test_job"]["data"]
1662 dagfile_expected = [
1663 f'JOB {job.name} "{job.subfile}"\n',
1664 f"DOT {self.dag.name}.dot\n",
1665 f"NODE_STATUS_FILE {self.dag.name}.node_status\n",
1666 ]
1668 self.dag.write(tmp_dir, "", "")
1670 self.assertIn("submit_path", self.dag.graph)
1671 self.assertEqual(self.dag.graph["submit_path"], tmp_dir)
1672 self.assertIn("dag_filename", self.dag.graph)
1673 self.assertEqual(self.dag.graph["dag_filename"], f"{self.dag.graph['name']}.dag")
1674 with open(os.path.join(tmp_dir, self.dag.graph["dag_filename"]), encoding="utf-8") as f:
1675 dagfile_actual = f.readlines()
1676 self.assertEqual(dagfile_actual, dagfile_expected)
1677 with open(os.path.join(tmp_dir, job.subfile), encoding="utf-8") as f:
1678 subfile_actual = f.readlines()
1679 self.assertEqual(subfile_actual, self.subfile_expected)
1682if __name__ == "__main__":
1683 unittest.main()