Coverage for python/lsst/ctrl/bps/htcondor/htcondor_service.py: 59%
262 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-02 05:18 -0400
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-02 05:18 -0400
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"""Interface between generic workflow to HTCondor workflow system."""
30__all__ = ["HTCondorService"]
33import logging
34import os
35from pathlib import Path
37import htcondor
38from packaging import version
40from lsst.ctrl.bps import (
41 BaseWmsService,
42 WmsStates,
43)
44from lsst.ctrl.bps.bps_utils import chdir
45from lsst.daf.butler import Config
46from lsst.utils.timer import time_this
48from .common_utils import WmsIdType, _wms_id_to_cluster, _wms_id_to_dir, _wms_id_type
49from .dagman_configurator import DagmanConfigurator
50from .htcondor_config import HTC_DEFAULTS_URI
51from .htcondor_workflow import HTCondorWorkflow
52from .lssthtc import (
53 _locate_schedds,
54 _update_rescue_file,
55 condor_q,
56 htc_backup_files,
57 htc_create_submit_from_cmd,
58 htc_create_submit_from_dag,
59 htc_create_submit_from_file,
60 htc_submit_dag,
61 htc_version,
62 read_dag_info,
63 read_dag_status,
64 write_dag_info,
65)
66from .provisioner import Provisioner
67from .report_utils import (
68 _get_status_from_id,
69 _get_status_from_path,
70 _report_from_id,
71 _report_from_path,
72 _summary_report,
73)
75_LOG = logging.getLogger(__name__)
78class HTCondorService(BaseWmsService):
79 """HTCondor version of WMS service."""
81 @property
82 def defaults(self):
83 return Config(HTC_DEFAULTS_URI)
85 @property
86 def defaults_uri(self):
87 return HTC_DEFAULTS_URI
89 def prepare(self, config, generic_workflow, out_prefix=None):
90 """Convert generic workflow to an HTCondor DAG ready for submission.
92 Parameters
93 ----------
94 config : `lsst.ctrl.bps.BpsConfig`
95 BPS configuration that includes necessary submit/runtime
96 information.
97 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
98 The generic workflow (e.g., has executable name and arguments).
99 out_prefix : `str`
100 The root directory into which all WMS-specific files are written.
102 Returns
103 -------
104 workflow : `lsst.ctrl.bps.htcondor.HTCondorWorkflow`
105 HTCondor workflow ready to be run.
106 """
107 _LOG.debug("out_prefix = '%s'", out_prefix)
108 with time_this(log=_LOG, level=logging.INFO, prefix=None, msg="Completed HTCondor workflow creation"):
109 _, enable_provisioning = config.search("provisionResources")
111 # If bps is doing provisioning, force a unique nodeset
112 # to reduce complications if user also manually does
113 # provisioning.
114 if enable_provisioning:
115 config[".bps_defined.nodeset"] = config[".bps_defined.timestamp"]
117 workflow = HTCondorWorkflow.from_generic_workflow(
118 config,
119 generic_workflow,
120 out_prefix,
121 f"{self.__class__.__module__}.{self.__class__.__name__}",
122 )
124 if enable_provisioning:
125 provisioner = Provisioner(config)
126 provisioner.configure()
127 provisioner.prepare("provisioningJob.bash", prefix=out_prefix)
128 provisioner.provision(workflow.dag)
130 try:
131 configurator = DagmanConfigurator(config)
132 except KeyError:
133 _LOG.debug(
134 "No DAGMan-specific settings were found in BPS config; "
135 "skipping writing DAG-specific configuration file."
136 )
137 else:
138 configurator.prepare("dagman.conf", prefix=out_prefix)
139 configurator.configure(workflow.dag)
141 with time_this(
142 log=_LOG, level=logging.INFO, prefix=None, msg="Completed writing out HTCondor workflow"
143 ):
144 workflow.write(out_prefix)
145 return workflow
147 def submit(self, workflow, **kwargs):
148 """Submit a single HTCondor workflow.
150 Parameters
151 ----------
152 workflow : `lsst.ctrl.bps.htcondor.HTCondorWorkflow`
153 A single HTCondor workflow to submit. run_id is updated after
154 successful submission to WMS.
155 **kwargs : `~typing.Any`
156 Keyword arguments for the options.
157 """
158 dag = workflow.dag
159 ver = version.parse(htc_version())
161 # For workflow portability, internal paths are all relative. Hence
162 # the DAG needs to be submitted to HTCondor from inside the submit
163 # directory.
164 with chdir(workflow.submit_path):
165 try:
166 if ver >= version.parse("8.9.3"): 166 ↛ 174line 166 didn't jump to line 174 because the condition on line 166 was always true
167 wms_config_path = None
168 if "bps_wms_config_path" in dag.graph["attr"]:
169 wms_config_path = dag.graph["attr"]["bps_wms_config_path"]
170 sub = htc_create_submit_from_dag(
171 dag.graph["dag_filename"], dag.graph["submit_options"], wms_config_path
172 )
173 else:
174 sub = htc_create_submit_from_cmd(dag.graph["dag_filename"], dag.graph["submit_options"])
175 except Exception:
176 _LOG.error(
177 "Problems creating HTCondor submit object from filename: %s", dag.graph["dag_filename"]
178 )
179 raise
181 _LOG.info("Submitting from directory: %s", os.getcwd())
182 schedd_dag_info = htc_submit_dag(sub)
183 if schedd_dag_info:
184 _, dag_info = next(iter(schedd_dag_info.items()))
185 dag_id, dag_ad = next(iter(dag_info.items()))
187 write_dag_info(f"{dag_ad['bps_run']}.info.json", schedd_dag_info)
189 dag.run_id = f"{dag_ad['ClusterId']}.{dag_ad['ProcId']}"
190 workflow.run_id = dag.run_id
191 else:
192 raise RuntimeError("Submission failed: unable to retrieve DAGMan job information")
194 def restart(self, wms_workflow_id):
195 """Restart a failed DAGMan workflow.
197 Parameters
198 ----------
199 wms_workflow_id : `str`
200 The directory with HTCondor files.
202 Returns
203 -------
204 run_id : `str`
205 HTCondor id of the restarted DAGMan job. If restart failed, it will
206 be set to None.
207 run_name : `str`
208 Name of the restarted workflow. If restart failed, it will be set
209 to None.
210 message : `str`
211 A message describing any issues encountered during the restart.
212 If there were no issues, an empty string is returned.
213 """
214 wms_path, id_type = _wms_id_to_dir(wms_workflow_id)
215 if wms_path is None:
216 return (
217 None,
218 None,
219 (
220 f"workflow with run id '{wms_workflow_id}' not found. "
221 "Hint: use run's submit directory as the id instead"
222 ),
223 )
225 if id_type in {WmsIdType.GLOBAL, WmsIdType.LOCAL}:
226 if not wms_path.is_dir(): 226 ↛ 229line 226 didn't jump to line 229 because the condition on line 226 was always true
227 return None, None, f"submit directory '{wms_path}' for run id '{wms_workflow_id}' not found."
229 _LOG.info("Restarting workflow from directory '%s'", wms_path)
230 rescue_dags = list(wms_path.glob("*.dag.rescue*"))
231 if not rescue_dags:
232 return None, None, f"HTCondor rescue DAG(s) not found in '{wms_path}'"
234 _LOG.info("Verifying that the workflow is not already in the job queue")
235 schedd_dag_info = condor_q(constraint=f'regexp("dagman$", Cmd) && Iwd == "{wms_path}"')
236 if schedd_dag_info:
237 _, dag_info = schedd_dag_info.popitem()
238 _, dag_ad = dag_info.popitem()
239 id_ = dag_ad["GlobalJobId"]
240 return None, None, f"Workflow already in the job queue (global job id: '{id_}')"
242 _LOG.info("Checking execution status of the workflow")
243 warn = False
244 dag_ad = read_dag_status(str(wms_path))
245 if dag_ad: 245 ↛ 254line 245 didn't jump to line 254 because the condition on line 245 was always true
246 nodes_total = dag_ad.get("NodesTotal", 0)
247 if nodes_total != 0: 247 ↛ 252line 247 didn't jump to line 252 because the condition on line 247 was always true
248 nodes_done = dag_ad.get("NodesDone", 0)
249 if nodes_total == nodes_done:
250 return None, None, "All jobs in the workflow finished successfully"
251 else:
252 warn = True
253 else:
254 warn = True
255 if warn: 255 ↛ 256line 255 didn't jump to line 256 because the condition on line 255 was never true
256 _LOG.warning(
257 "Cannot determine the execution status of the workflow, continuing with restart regardless"
258 )
260 # In the case of lazy DAGs, workflow summaries can change at
261 # runtime. So read the workflow's info.json file before moving
262 # it to backup dir and use to update the summaries later before
263 # writing the new info.json file.
264 dag_info_filename, old_dag_schedd_info = read_dag_info(wms_path)
265 old_dag_info = next(iter(old_dag_schedd_info.values()))
266 old_dag_ad = next(iter(old_dag_info.values()))
268 _LOG.info("Backing up select HTCondor files from previous run attempt")
269 rescue_files = sorted(wms_path.glob("*.rescue[0-9][0-9][0-9]"))
270 last_rescue_file = Path(rescue_files[-1]) if rescue_files else None
271 has_subdags = (wms_path / "subdags").exists()
272 failed_subdags = None
273 if last_rescue_file and has_subdags: 273 ↛ 274line 273 didn't jump to line 274 because the condition on line 273 was never true
274 failed_subdags = set(_update_rescue_file(last_rescue_file))
275 htc_backup_files(wms_path, subdir="backups", failed_subdags=failed_subdags)
277 # For workflow portability, internal paths are all relative. Hence
278 # the DAG needs to be resubmitted to HTCondor from inside the submit
279 # directory.
280 _LOG.info("Adding workflow to the job queue")
281 run_id, run_name, message = None, None, ""
282 with chdir(wms_path):
283 try:
284 dag_path = next(Path.cwd().glob("*.dag.condor.sub"))
285 except StopIteration:
286 message = f"DAGMan submit description file not found in '{wms_path}'"
287 else:
288 sub = htc_create_submit_from_file(dag_path.name)
289 schedd_dag_info = htc_submit_dag(sub)
291 # Save select information about the DAGMan job to a file. Use
292 # the run name (available in the ClassAd) as the filename.
293 if schedd_dag_info:
294 dag_info = next(iter(schedd_dag_info.values()))
295 dag_ad = next(iter(dag_info.values()))
297 # Just in case lazy DAGs, update the summaries.
298 dag_ad["bps_job_summary"] = old_dag_ad["bps_job_summary"]
299 dag_ad["bps_run_quanta"] = old_dag_ad["bps_run_quanta"]
301 write_dag_info(dag_info_filename, schedd_dag_info)
302 run_id = f"{dag_ad['ClusterId']}.{dag_ad['ProcId']}"
303 run_name = dag_ad["bps_run"]
304 else:
305 message = "DAGMan job information unavailable"
307 return run_id, run_name, message
309 def list_submitted_jobs(self, wms_id=None, user=None, require_bps=True, pass_thru=None, is_global=False):
310 """Query WMS for list of submitted WMS workflows/jobs.
312 This should be a quick lookup function to create list of jobs for
313 other functions.
315 Parameters
316 ----------
317 wms_id : `int` or `str`, optional
318 Id or path that can be used by WMS service to look up job.
319 user : `str`, optional
320 User whose submitted jobs should be listed.
321 require_bps : `bool`, optional
322 Whether to require jobs returned in list to be bps-submitted jobs.
323 pass_thru : `str`, optional
324 Information to pass through to WMS.
325 is_global : `bool`, optional
326 If set, all job queues (and their histories) will be queried for
327 job information. Defaults to False which means that only the local
328 job queue will be queried.
330 Returns
331 -------
332 job_ids : `list` [`~typing.Any`]
333 Only job ids to be used by cancel and other functions. Typically
334 this means top-level jobs (i.e., not children jobs).
335 """
336 _LOG.debug(
337 "list_submitted_jobs params: wms_id=%s, user=%s, require_bps=%s, pass_thru=%s, is_global=%s",
338 wms_id,
339 user,
340 require_bps,
341 pass_thru,
342 is_global,
343 )
345 # Determine which Schedds will be queried for job information.
346 coll = htcondor.Collector()
348 schedd_ads = []
349 if is_global:
350 schedd_ads.extend(coll.locateAll(htcondor.DaemonTypes.Schedd))
351 else:
352 schedd_ads.append(coll.locate(htcondor.DaemonTypes.Schedd))
354 # Construct appropriate constraint expression using provided arguments.
355 constraint = "False"
356 if wms_id is None:
357 if user is not None:
358 constraint = f'(Owner == "{user}")'
359 else:
360 schedd_ad, cluster_id, id_type = _wms_id_to_cluster(wms_id)
361 if cluster_id is not None:
362 constraint = f"(DAGManJobId == {cluster_id} || ClusterId == {cluster_id})"
364 # If provided id is either a submission path or a global id,
365 # make sure the right Schedd will be queried regardless of
366 # 'is_global' value.
367 if id_type in {WmsIdType.GLOBAL, WmsIdType.PATH}:
368 schedd_ads = [schedd_ad]
369 if require_bps:
370 constraint += ' && (bps_isjob == "True")'
371 if pass_thru:
372 if "-forcex" in pass_thru:
373 pass_thru_2 = pass_thru.replace("-forcex", "")
374 if pass_thru_2 and not pass_thru_2.isspace():
375 constraint += f" && ({pass_thru_2})"
376 else:
377 constraint += f" && ({pass_thru})"
379 # Create a list of scheduler daemons which need to be queried.
380 schedds = {ad["Name"]: htcondor.Schedd(ad) for ad in schedd_ads}
382 _LOG.debug("constraint = %s, schedds = %s", constraint, ", ".join(schedds))
383 results = condor_q(constraint=constraint, schedds=schedds)
385 # Prune child jobs where DAG job is in queue (i.e., aren't orphans).
386 job_ids = []
387 for job_info in results.values():
388 for job_id, job_ad in job_info.items():
389 _LOG.debug("job_id=%s DAGManJobId=%s", job_id, job_ad.get("DAGManJobId", "None"))
390 if "DAGManJobId" not in job_ad:
391 job_ids.append(job_ad.get("GlobalJobId", job_id))
392 else:
393 _LOG.debug("Looking for %s", f"{job_ad['DAGManJobId']}.0")
394 _LOG.debug("\tin jobs.keys() = %s", job_info.keys())
395 if f"{job_ad['DAGManJobId']}.0" not in job_info: # orphaned job
396 job_ids.append(job_ad.get("GlobalJobId", job_id))
398 _LOG.debug("job_ids = %s", job_ids)
399 return job_ids
401 def get_status(
402 self,
403 wms_workflow_id: str,
404 hist: float = 1,
405 is_global: bool = False,
406 ) -> tuple[WmsStates, str]:
407 """Return status of run based upon given constraints.
409 Parameters
410 ----------
411 wms_workflow_id : `str`
412 Limit to specific run based on id (queue id or path).
413 hist : `float`, optional
414 Limit history search to this many days. Defaults to 1.
415 is_global : `bool`, optional
416 If set, all job queues (and their histories) will be queried for
417 job information. Defaults to False which means that only the local
418 job queue will be queried.
420 Returns
421 -------
422 state : `lsst.ctrl.bps.WmsStates`
423 Status of single run from given information.
424 message : `str`
425 Extra message for status command to print. This could be pointers
426 to documentation or to WMS specific commands.
427 """
428 _LOG.debug("get_status: id=%s, hist=%s, is_global=%s", wms_workflow_id, hist, is_global)
430 id_type = _wms_id_type(wms_workflow_id)
431 _LOG.debug("id_type = %s", id_type.name)
433 if id_type == WmsIdType.LOCAL:
434 schedulers = _locate_schedds(locate_all=is_global)
435 _LOG.debug("schedulers = %s", schedulers)
436 state, message = _get_status_from_id(wms_workflow_id, hist, schedds=schedulers)
437 elif id_type == WmsIdType.GLOBAL:
438 schedulers = _locate_schedds(locate_all=True)
439 _LOG.debug("schedulers = %s", schedulers)
440 state, message = _get_status_from_id(wms_workflow_id, hist, schedds=schedulers)
441 elif id_type == WmsIdType.PATH:
442 state, message = _get_status_from_path(wms_workflow_id)
443 else:
444 state, message = WmsStates.UNKNOWN, "Invalid job id"
445 _LOG.debug("state: %s, %s", state, message)
447 return state, message
449 def report(
450 self,
451 wms_workflow_id=None,
452 user=None,
453 hist=0,
454 pass_thru=None,
455 is_global=False,
456 return_exit_codes=False,
457 ):
458 """Return run information based upon given constraints.
460 Parameters
461 ----------
462 wms_workflow_id : `str`, optional
463 Limit to specific run based on id.
464 user : `str`, optional
465 Limit results to runs for this user.
466 hist : `float`, optional
467 Limit history search to this many days. Defaults to 0.
468 pass_thru : `str`, optional
469 Constraints to pass through to HTCondor.
470 is_global : `bool`, optional
471 If set, all job queues (and their histories) will be queried for
472 job information. Defaults to False which means that only the local
473 job queue will be queried.
474 return_exit_codes : `bool`, optional
475 If set, return exit codes related to jobs with a
476 non-success status. Defaults to False, which means that only
477 the summary state is returned.
479 Only applicable in the context of a WMS with associated
480 handlers to return exit codes from jobs.
482 Returns
483 -------
484 runs : `list` [`lsst.ctrl.bps.WmsRunReport`]
485 Information about runs from given job information.
486 message : `str`
487 Extra message for report command to print. This could be pointers
488 to documentation or to WMS specific commands.
489 """
490 if wms_workflow_id:
491 id_type = _wms_id_type(wms_workflow_id)
492 if id_type == WmsIdType.LOCAL:
493 schedulers = _locate_schedds(locate_all=is_global)
494 run_reports, message = _report_from_id(wms_workflow_id, hist, schedds=schedulers)
495 elif id_type == WmsIdType.GLOBAL:
496 schedulers = _locate_schedds(locate_all=True)
497 run_reports, message = _report_from_id(wms_workflow_id, hist, schedds=schedulers)
498 elif id_type == WmsIdType.PATH:
499 run_reports, message = _report_from_path(wms_workflow_id)
500 else:
501 run_reports, message = {}, "Invalid job id"
502 else:
503 schedulers = _locate_schedds(locate_all=is_global)
504 run_reports, message = _summary_report(user, hist, pass_thru, schedds=schedulers)
505 _LOG.debug("report: %s, %s", run_reports, message)
507 return list(run_reports.values()), message
509 def cancel(self, wms_id, pass_thru=None):
510 """Cancel submitted workflows/jobs.
512 Parameters
513 ----------
514 wms_id : `str`
515 Id or path of job that should be canceled.
516 pass_thru : `str`, optional
517 Information to pass through to WMS.
519 Returns
520 -------
521 deleted : `bool`
522 Whether successful deletion or not. Currently, if any doubt or any
523 individual jobs not deleted, return False.
524 message : `str`
525 Any message from WMS (e.g., error details).
526 """
527 _LOG.debug("Canceling wms_id = %s", wms_id)
529 schedd_ad, cluster_id, _ = _wms_id_to_cluster(wms_id)
531 if cluster_id is None:
532 deleted = False
533 message = "invalid id"
534 else:
535 _LOG.debug(
536 "Canceling job managed by schedd_name = %s with cluster_id = %s",
537 cluster_id,
538 schedd_ad["Name"],
539 )
540 schedd = htcondor.Schedd(schedd_ad)
542 constraint = f"ClusterId == {cluster_id}"
543 if pass_thru is not None and "-forcex" in pass_thru:
544 pass_thru_2 = pass_thru.replace("-forcex", "")
545 if pass_thru_2 and not pass_thru_2.isspace():
546 constraint += f"&& ({pass_thru_2})"
547 _LOG.debug("JobAction.RemoveX constraint = %s", constraint)
548 results = schedd.act(htcondor.JobAction.RemoveX, constraint)
549 else:
550 if pass_thru:
551 constraint += f"&& ({pass_thru})"
552 _LOG.debug("JobAction.Remove constraint = %s", constraint)
553 results = schedd.act(htcondor.JobAction.Remove, constraint)
554 _LOG.debug("Remove results: %s", results)
556 if results["TotalSuccess"] > 0 and results["TotalError"] == 0:
557 deleted = True
558 message = ""
559 else:
560 deleted = False
561 if results["TotalSuccess"] == 0 and results["TotalError"] == 0:
562 message = "no such bps job in batch queue"
563 else:
564 message = f"unknown problems deleting: {results}"
566 _LOG.debug("deleted: %s; message = %s", deleted, message)
567 return deleted, message
569 def ping(self, pass_thru):
570 """Check whether WMS services are up, reachable, and can authenticate
571 if authentication is required.
573 The services to be checked are those needed for submit, report, cancel,
574 restart, but ping cannot guarantee whether jobs would actually run
575 successfully.
577 Parameters
578 ----------
579 pass_thru : `str`, optional
580 Information to pass through to WMS.
582 Returns
583 -------
584 status : `int`
585 0 for success, non-zero for failure.
586 message : `str`
587 Any message from WMS (e.g., error details).
588 """
589 coll = htcondor.Collector()
590 secman = htcondor.SecMan()
591 status = 0
592 message = ""
593 _LOG.info("Not verifying that compute resources exist.")
594 try:
595 for daemon_type in [htcondor.DaemonTypes.Schedd, htcondor.DaemonTypes.Collector]:
596 _ = secman.ping(coll.locate(daemon_type))
597 except htcondor.HTCondorLocateError:
598 status = 1
599 message = f"Could not locate {daemon_type} service."
600 except htcondor.HTCondorIOError:
601 status = 1
602 message = f"Permission problem with {daemon_type} service."
603 return status, message
605 def run_submission_checks(self):
606 """Check to run at start if running WMS specific submission steps.
608 Any exception other than NotImplementedError will halt submission.
609 Submit directory may not yet exist when this is called.
610 """
611 # Some early config sanity checks
612 found, value = self.config.search("bpsMakeCommand")
613 bps_make_command = value if found else True
614 if not bps_make_command:
615 found, value = self.config.search("payloadCommand", opt={"replaceVars": False})
616 if not found:
617 raise KeyError("Missing 'payloadCommand' in config while bpsMakeCommand=True")
619 if "setupEnv" in value:
620 found, value = self.config.search("setupEnv", opt={"replaceVars": False})
621 if not found:
622 raise KeyError("Missing 'setupEnv' in config, but appears in payloadCommand")
624 if "lsstVersion" in value:
625 found, value = self.config.search("lsstVersion", opt={"replaceVars": False})
626 if not found:
627 raise KeyError("Missing 'lsstVersion' in config, but appears in setupEnv in config")