Coverage for tests/test_common_utils.py: 100%
131 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-16 02:16 -0700
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-16 02:16 -0700
1# This file is part of ctrl_bps_htcondor.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This software is dual licensed under the GNU General Public License and also
10# under a 3-clause BSD license. Recipients may choose which of these licenses
11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt,
12# respectively. If you choose the GPL option then the following text applies
13# (but note that there is still no warranty even if you opt for BSD instead):
14#
15# This program is free software: you can redistribute it and/or modify
16# it under the terms of the GNU General Public License as published by
17# the Free Software Foundation, either version 3 of the License, or
18# (at your option) any later version.
19#
20# This program is distributed in the hope that it will be useful,
21# but WITHOUT ANY WARRANTY; without even the implied warranty of
22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
23# GNU General Public License for more details.
24#
25# You should have received a copy of the GNU General Public License
26# along with this program. If not, see <https://www.gnu.org/licenses/>.
28"""Unit tests for the common utility functions."""
30import logging
31import os
32import unittest
33from pathlib import Path
35import htcondor
37from lsst.ctrl.bps import (
38 WmsStates,
39)
40from lsst.ctrl.bps.htcondor import common_utils
41from lsst.utils.tests import temporaryDirectory
43logger = logging.getLogger("lsst.ctrl.bps.htcondor")
46class HtcNodeStatusToWmsStateTestCase(unittest.TestCase):
47 """Test assigning WMS state base on HTCondor node status."""
49 def setUp(self):
50 pass
52 def tearDown(self):
53 pass
55 def testNotReady(self):
56 job = {"NodeStatus": common_utils.NodeStatus.NOT_READY}
57 result = common_utils._htc_node_status_to_wms_state(job)
58 self.assertEqual(result, WmsStates.UNREADY)
60 def testReady(self):
61 job = {"NodeStatus": common_utils.NodeStatus.READY}
62 result = common_utils._htc_node_status_to_wms_state(job)
63 self.assertEqual(result, WmsStates.READY)
65 def testPrerun(self):
66 job = {"NodeStatus": common_utils.NodeStatus.PRERUN}
67 result = common_utils._htc_node_status_to_wms_state(job)
68 self.assertEqual(result, WmsStates.MISFIT)
70 def testSubmittedHeld(self):
71 job = {
72 "NodeStatus": common_utils.NodeStatus.SUBMITTED,
73 "JobProcsHeld": 1,
74 "StatusDetails": "",
75 "JobProcsQueued": 0,
76 }
77 result = common_utils._htc_node_status_to_wms_state(job)
78 self.assertEqual(result, WmsStates.HELD)
80 def testSubmittedRunning(self):
81 job = {
82 "NodeStatus": common_utils.NodeStatus.SUBMITTED,
83 "JobProcsHeld": 0,
84 "StatusDetails": "not_idle",
85 "JobProcsQueued": 0,
86 }
87 result = common_utils._htc_node_status_to_wms_state(job)
88 self.assertEqual(result, WmsStates.RUNNING)
90 def testSubmittedPending(self):
91 job = {
92 "NodeStatus": common_utils.NodeStatus.SUBMITTED,
93 "JobProcsHeld": 0,
94 "StatusDetails": "",
95 "JobProcsQueued": 1,
96 }
97 result = common_utils._htc_node_status_to_wms_state(job)
98 self.assertEqual(result, WmsStates.PENDING)
100 def testPostrun(self):
101 job = {"NodeStatus": common_utils.NodeStatus.POSTRUN}
102 result = common_utils._htc_node_status_to_wms_state(job)
103 self.assertEqual(result, WmsStates.MISFIT)
105 def testDone(self):
106 job = {"NodeStatus": common_utils.NodeStatus.DONE}
107 result = common_utils._htc_node_status_to_wms_state(job)
108 self.assertEqual(result, WmsStates.SUCCEEDED)
110 def testErrorDagmanSuccess(self):
111 job = {
112 "NodeStatus": common_utils.NodeStatus.ERROR,
113 "StatusDetails": "DAGMAN error 0",
114 }
115 result = common_utils._htc_node_status_to_wms_state(job)
116 self.assertEqual(result, WmsStates.SUCCEEDED)
118 def testErrorDagmanFailure(self):
119 job = {
120 "NodeStatus": common_utils.NodeStatus.ERROR,
121 "StatusDetails": "DAGMAN error 1",
122 }
123 result = common_utils._htc_node_status_to_wms_state(job)
124 self.assertEqual(result, WmsStates.FAILED)
126 def testFutile(self):
127 job = {"NodeStatus": common_utils.NodeStatus.FUTILE}
128 result = common_utils._htc_node_status_to_wms_state(job)
129 self.assertEqual(result, WmsStates.PRUNED)
131 def testDeletedJob(self):
132 job = {
133 "NodeStatus": common_utils.NodeStatus.ERROR,
134 "StatusDetails": "HTCondor reported ULOG_JOB_ABORTED event for job proc (1.0.0)",
135 "JobProcsQueued": 0,
136 }
137 result = common_utils._htc_node_status_to_wms_state(job)
138 self.assertEqual(result, WmsStates.DELETED)
141class HtcStatusToWmsStateTestCase(unittest.TestCase):
142 """Test assigning WMS state base on HTCondor status."""
144 def testJobStatus(self):
145 job = {
146 "ClusterId": 1,
147 "JobStatus": htcondor.JobStatus.IDLE,
148 "bps_job_label": "foo",
149 }
150 result = common_utils._htc_status_to_wms_state(job)
151 self.assertEqual(result, WmsStates.PENDING)
153 def testNodeStatus(self):
154 # Hold/Release test case
155 job = {
156 "ClusterId": 1,
157 "JobStatus": None,
158 "NodeStatus": common_utils.NodeStatus.SUBMITTED,
159 "JobProcsHeld": 0,
160 "StatusDetails": "",
161 "JobProcsQueued": 1,
162 }
163 result = common_utils._htc_status_to_wms_state(job)
164 self.assertEqual(result, WmsStates.PENDING)
166 def testNeitherStatus(self):
167 job = {"ClusterId": 1}
168 result = common_utils._htc_status_to_wms_state(job)
169 self.assertEqual(result, WmsStates.MISFIT)
171 def testRetrySuccess(self):
172 job = {
173 "NodeStatus": 5,
174 "Node": "8e62c569-ae2e-44e8-be36-d1aee333a129_isr_903342_10",
175 "RetryCount": 0,
176 "ClusterId": 851,
177 "ProcId": 0,
178 "MyType": "JobTerminatedEvent",
179 "EventTypeNumber": 5,
180 "HoldReasonCode": 3,
181 "HoldReason": "Job raised a signal 9. Handling signal as if job has gone over memory limit.",
182 "HoldReasonSubCode": 34,
183 "ToE": {
184 "ExitBySignal": False,
185 "ExitCode": 0,
186 },
187 "JobStatus": htcondor.JobStatus.COMPLETED,
188 "ExitBySignal": False,
189 "ExitCode": 0,
190 }
191 result = common_utils._htc_status_to_wms_state(job)
192 self.assertEqual(result, WmsStates.SUCCEEDED)
195class WmsIdToDirTestCase(unittest.TestCase):
196 """Test _wms_id_to_dir function."""
198 @unittest.mock.patch("lsst.ctrl.bps.htcondor.common_utils._wms_id_type")
199 def testInvalidIdType(self, _wms_id_type_mock):
200 _wms_id_type_mock.return_value = common_utils.WmsIdType.UNKNOWN
201 with self.assertRaises(TypeError) as cm:
202 _, _ = common_utils._wms_id_to_dir("not_used")
203 self.assertIn("Invalid job id type", str(cm.exception))
205 @unittest.mock.patch("lsst.ctrl.bps.htcondor.common_utils._wms_id_type")
206 def testAbsPathId(self, mock_wms_id_type):
207 mock_wms_id_type.return_value = common_utils.WmsIdType.PATH
208 with temporaryDirectory() as tmp_dir:
209 wms_path, id_type = common_utils._wms_id_to_dir(tmp_dir)
210 self.assertEqual(id_type, common_utils.WmsIdType.PATH)
211 self.assertEqual(Path(tmp_dir).resolve(), wms_path)
213 @unittest.mock.patch("lsst.ctrl.bps.htcondor.common_utils._wms_id_type")
214 def testRelPathId(self, _wms_id_type_mock):
215 _wms_id_type_mock.return_value = common_utils.WmsIdType.PATH
216 orig_dir = Path.cwd()
217 with temporaryDirectory() as tmp_dir:
218 os.chdir(tmp_dir)
219 abs_path = Path(tmp_dir) / "newdir"
220 abs_path.mkdir()
221 wms_path, id_type = common_utils._wms_id_to_dir("newdir")
222 self.assertEqual(id_type, common_utils.WmsIdType.PATH)
223 self.assertEqual(abs_path.resolve(), wms_path)
224 os.chdir(orig_dir)
227class WmsIdTypeTestCase(unittest.TestCase):
228 """Test _wms_id_type function."""
230 def testIntId(self):
231 id_type = common_utils._wms_id_type("4")
232 self.assertEqual(id_type, common_utils.WmsIdType.LOCAL)
234 def testPathId(self):
235 with temporaryDirectory() as tmp_dir:
236 id_type = common_utils._wms_id_type(str(tmp_dir))
237 self.assertEqual(id_type, common_utils.WmsIdType.PATH)
239 def testGlobalId(self):
240 id_type = common_utils._wms_id_type("testmachine#5044.0#1757720957")
241 self.assertEqual(id_type, common_utils.WmsIdType.GLOBAL)
243 def testUnknownType(self):
244 id_type = common_utils._wms_id_type(["bad param"])
245 self.assertEqual(id_type, common_utils.WmsIdType.UNKNOWN)
248class WmsIdToClusterTestCase(unittest.TestCase):
249 """Test _wms_id_to_cluster function."""
251 class _MockCollector:
252 def locate(self, daemon_type, schedd_name):
253 return """[
254 CondorPlatform = "$CondorPlatform: X86_64-AlmaLinux_9.6 $";
255 MyType = "Scheduler";
256 Machine = "testmachine";
257 Name = "testmachine";
258 CondorVersion = "$CondorVersion: 24.0.10 2025-08-05 $";
259 MyAddress = "<127.0.0.1:9618?addrs=127.0.0.1-9618+snip>"
260 ]"""
262 @unittest.mock.patch("htcondor.Collector", new=_MockCollector)
263 @unittest.mock.patch("lsst.ctrl.bps.htcondor.common_utils.read_dag_info")
264 def testPath(self, mock_read):
265 # path must exist or _wms_id_type will assume GLOBAL string
266 with temporaryDirectory() as tmp_dir:
267 mock_read.return_value = [
268 "dummy_str",
269 {
270 "testmachine": {
271 "1163.0": {
272 "testmachine": {
273 "ClusterId": 1163,
274 "GlobalJobId": "testmachine#1163.0#1722040518",
275 }
276 }
277 }
278 },
279 ]
280 schedd_ad, cluster_id, id_type = common_utils._wms_id_to_cluster(tmp_dir)
281 self.assertEqual(id_type, common_utils.WmsIdType.PATH)
282 self.assertEqual(cluster_id, 1163)
285if __name__ == "__main__":
286 unittest.main()