Coverage for python/lsst/ctrl/bps/htcondor/lssthtc.py: 79%
954 statements
« prev ^ index » next coverage.py v7.16.2, created at 2026-09-30 11:10 +0000
« prev ^ index » next coverage.py v7.16.2, created at 2026-09-30 11:10 +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/>.
28"""Placeholder HTCondor DAGMan API.
30There is new work on a python DAGMan API from HTCondor. However, at this
31time, it tries to make things easier by assuming DAG is easily broken into
32levels where there are 1-1 or all-to-all relationships to nodes in next
33level. LSST workflows are more complicated.
34"""
36__all__ = [
37 "MISSING_ID",
38 "DagStatus",
39 "HTCDag",
40 "HTCJob",
41 "NodeStatus",
42 "RestrictedDict",
43 "WmsNodeType",
44 "condor_history",
45 "condor_q",
46 "condor_search",
47 "condor_status",
48 "htc_backup_files",
49 "htc_check_dagman_output",
50 "htc_create_submit_from_cmd",
51 "htc_create_submit_from_dag",
52 "htc_create_submit_from_file",
53 "htc_escape",
54 "htc_query_history",
55 "htc_query_present",
56 "htc_submit_dag",
57 "htc_tweak_log_info",
58 "htc_version",
59 "htc_write_attribs",
60 "htc_write_condor_file",
61 "pegasus_name_to_label",
62 "read_dag_info",
63 "read_dag_log",
64 "read_dag_nodes_log",
65 "read_dag_status",
66 "read_node_status",
67 "summarize_dag",
68 "update_job_info",
69 "write_dag_info",
70]
73import itertools
74import json
75import logging
76import os
77import pprint
78import re
79import subprocess
80from collections import Counter, defaultdict
81from collections.abc import MutableMapping
82from datetime import datetime, timedelta
83from enum import IntEnum, auto
84from pathlib import Path
85from typing import Any, TextIO
87import classad
88import htcondor
89import networkx
90from deprecated.sphinx import deprecated
91from packaging import version
93from .handlers import HTC_JOB_AD_HANDLERS
95_LOG = logging.getLogger(__name__)
97MISSING_ID = "-99999"
100class DagStatus(IntEnum):
101 """HTCondor DAGMan's statuses for a DAG."""
103 OK = 0
104 ERROR = 1 # an error condition different than those listed here
105 FAILED = 2 # one or more nodes in the DAG have failed
106 ABORTED = 3 # the DAG has been aborted by an ABORT-DAG-ON specification
107 REMOVED = 4 # the DAG has been removed by condor_rm
108 CYCLE = 5 # a cycle was found in the DAG
109 SUSPENDED = 6 # the DAG has been suspended (see section 2.10.8)
112@deprecated(
113 reason="The JobStatus is internally replaced by htcondor.JobStatus. "
114 "External reporting code should be using ctrl_bps.WmsStates. "
115 "This class will be removed after v30.",
116 version="v30.0",
117 category=FutureWarning,
118)
119class JobStatus(IntEnum):
120 """HTCondor's statuses for jobs."""
122 UNEXPANDED = 0 # Unexpanded
123 IDLE = 1 # Idle
124 RUNNING = 2 # Running
125 REMOVED = 3 # Removed
126 COMPLETED = 4 # Completed
127 HELD = 5 # Held
128 TRANSFERRING_OUTPUT = 6 # Transferring_Output
129 SUSPENDED = 7 # Suspended
132class NodeStatus(IntEnum):
133 """HTCondor's statuses for DAGman nodes."""
135 # (STATUS_NOT_READY): At least one parent has not yet finished or the node
136 # is a FINAL node.
137 NOT_READY = 0
139 # (STATUS_READY): All parents have finished, but the node is not yet
140 # running.
141 READY = 1
143 # (STATUS_PRERUN): The node’s PRE script is running.
144 PRERUN = 2
146 # (STATUS_SUBMITTED): The node’s HTCondor job(s) are in the queue.
147 # StatusDetails = "not_idle" -> running.
148 # JobProcsHeld = 1-> hold.
149 # JobProcsQueued = 1 -> idle.
150 SUBMITTED = 3
152 # (STATUS_POSTRUN): The node’s POST script is running.
153 POSTRUN = 4
155 # (STATUS_DONE): The node has completed successfully.
156 DONE = 5
158 # (STATUS_ERROR): The node has failed. StatusDetails has info (e.g.,
159 # ULOG_JOB_ABORTED for deleted job).
160 ERROR = 6
162 # (STATUS_FUTILE): The node will never run because ancestor node failed.
163 FUTILE = 7
166class WmsNodeType(IntEnum):
167 """HTCondor plugin node types to help with payload reporting."""
169 UNKNOWN = auto()
170 """Dummy value when missing."""
172 PAYLOAD = auto()
173 """Payload job."""
175 FINAL = auto()
176 """Final job."""
178 SERVICE = auto()
179 """Service job."""
181 NOOP = auto()
182 """NOOP job used for ordering jobs."""
184 SUBDAG = auto()
185 """SUBDAG job used for ordering jobs."""
187 SUBDAG_CHECK = auto()
188 """Job used to correctly prune jobs after a subdag."""
191HTC_QUOTE_KEYS = {"environment", "arguments"}
192HTC_VALID_JOB_KEYS = {
193 "universe",
194 "executable",
195 "arguments",
196 "environment",
197 "log",
198 "error",
199 "output",
200 "should_transfer_files",
201 "when_to_transfer_output",
202 "getenv",
203 "notification",
204 "notify_user",
205 "concurrency_limit",
206 "transfer_executable",
207 "transfer_input_files",
208 "transfer_output_files",
209 "transfer_output_remaps",
210 "request_cpus",
211 "request_memory",
212 "request_disk",
213 "priority",
214 "category",
215 "requirements",
216 "on_exit_hold",
217 "on_exit_hold_reason",
218 "on_exit_hold_subcode",
219 "max_retries",
220 "retry_until",
221 "periodic_release",
222 "periodic_remove",
223 "accounting_group",
224 "accounting_group_user",
225 "kill_sig",
226 "want_graceful_removal",
227 "job_max_vacate_time",
228}
229HTC_VALID_JOB_DAG_KEYS = {
230 "dir",
231 "noop",
232 "done",
233 "vars",
234 "pre",
235 "post",
236 "retry",
237 "retry_unless_exit",
238 "abort_dag_on",
239 "abort_exit",
240 "priority",
241}
242HTC_VERSION = version.parse(htcondor.__version__)
245class RestrictedDict(MutableMapping):
246 """A dictionary that only allows certain keys.
248 Parameters
249 ----------
250 valid_keys : `~collections.abc.Container`
251 Strings that are valid keys.
252 init_data : `dict` or `RestrictedDict`, optional
253 Initial data.
255 Raises
256 ------
257 KeyError
258 If invalid key(s) in init_data.
259 """
261 def __init__(self, valid_keys, init_data=()):
262 self.valid_keys = valid_keys
263 self.data = {}
264 self.update(init_data)
266 def __getitem__(self, key):
267 """Return value for given key if exists.
269 Parameters
270 ----------
271 key : `str`
272 Identifier for value to return.
274 Returns
275 -------
276 value : `~typing.Any`
277 Value associated with given key.
279 Raises
280 ------
281 KeyError
282 If key doesn't exist.
283 """
284 return self.data[key]
286 def __delitem__(self, key):
287 """Delete value for given key if exists.
289 Parameters
290 ----------
291 key : `str`
292 Identifier for value to delete.
294 Raises
295 ------
296 KeyError
297 If key doesn't exist.
298 """
299 del self.data[key]
301 def __setitem__(self, key, value):
302 """Store key,value in internal dict only if key is valid.
304 Parameters
305 ----------
306 key : `str`
307 Identifier to associate with given value.
308 value : `~typing.Any`
309 Value to store.
311 Raises
312 ------
313 KeyError
314 If key is invalid.
315 """
316 if key not in self.valid_keys: 316 ↛ 317line 316 didn't jump to line 317 because the condition on line 316 was never true
317 raise KeyError(f"Invalid key {key}")
318 self.data[key] = value
320 def __iter__(self):
321 return self.data.__iter__()
323 def __len__(self):
324 return len(self.data)
326 def __str__(self):
327 return str(self.data)
330def htc_backup_files(
331 wms_path: str | os.PathLike,
332 subdir: str | os.PathLike | None = None,
333 limit: int = 100,
334 failed_subdags: set[str] | None = None,
335) -> None:
336 """Backup select HTCondor files in the submit directory.
338 Files will be saved in separate subdirectories which will be created in
339 the submit directory where the files are located. These subdirectories
340 will be consecutive, zero-padded integers. Their values will correspond to
341 the number of HTCondor rescue DAGs in the submit directory.
343 Hence, with the default settings, copies after the initial failed run will
344 be placed in '001' subdirectory, '002' after the first restart, and so on
345 until the limit of backups is reached. If there's no rescue DAG yet, files
346 will be copied to '000' subdirectory.
348 This is not a generic function for making backups. It is intended to be
349 used once, just before a restart, to make snapshots of files which will be
350 overwritten by HTCondor after during the next run.
352 Parameters
353 ----------
354 wms_path : `str` or `os.PathLike`
355 Path to the submit directory either absolute or relative.
356 subdir : `str` or `os.PathLike`, optional
357 A path, relative to the submit directory, where all subdirectories with
358 backup files will be kept. Defaults to None which means that the backup
359 subdirectories will be placed directly in the submit directory.
360 limit : `int`, optional
361 Maximal number of backups. If the number of backups reaches the limit,
362 the last backup files will be overwritten. The default value is 100
363 to match the default value of HTCondor's DAGMAN_MAX_RESCUE_NUM in
364 version 8.8+.
365 failed_subdags : `set` [`str`], optional
366 Names of subdag jobs that failed. Only files of these subdags will be
367 backed up. If None (the default), files of all subdags with rescue
368 files are backed up.
370 Raises
371 ------
372 FileNotFoundError
373 If the submit directory or the file that needs to be backed up does not
374 exist.
375 OSError
376 If the submit directory cannot be accessed or backing up a file failed
377 either due to permission or filesystem related issues.
378 """
379 width = len(str(limit))
381 path = Path(wms_path).resolve()
382 if not path.is_dir():
383 raise FileNotFoundError(f"Directory {path} not found")
385 # Initialize the backup counter.
386 # If using control DAG, don't want to include nested DAGs.
387 rescue_dags = list(path.glob("*_ctrl.dag.rescue[0-9][0-9][0-9]"))
388 if not rescue_dags: 388 ↛ 390line 388 didn't jump to line 390 because the condition on line 388 was always true
389 rescue_dags = list(path.glob("*.rescue[0-9][0-9][0-9]"))
390 counter = min(len(rescue_dags), limit)
392 # Create the backup directory and move select files there.
393 dest = path
394 if subdir:
395 # PurePath.is_relative_to() is not available before Python 3.9. Hence
396 # we need to check is 'subdir' is in the submit directory in some other
397 # way if it is an absolute path.
398 subdir = Path(subdir)
399 if subdir.is_absolute():
400 subdir = subdir.resolve() # Since resolve was run on path, must run it here
401 if dest not in subdir.parents:
402 _LOG.warning(
403 "Invalid backup location: '%s' not in the submit directory, will use '%s' instead.",
404 subdir,
405 wms_path,
406 )
407 else:
408 dest /= subdir
409 else:
410 dest /= subdir
411 dest /= f"{counter:0{width}}"
412 _LOG.debug("dest = %s", dest)
413 try:
414 dest.mkdir(parents=True, exist_ok=False if counter < limit else True)
415 except FileExistsError:
416 _LOG.warning("Refusing to do backups: target directory '%s' already exists", dest)
417 else:
418 htc_backup_files_single_path(path, dest)
420 # Back up selected files for failed subdags as well.
421 #
422 # Do NOT back up files of subdags that succeeded! These files need to stay
423 # in their respective directories as they will not be recreated by HTCondor
424 # after the run is restarted. As HTCondorService.report() uses information
425 # in these files to determine job statuses, their absence may lead to
426 # reporting incorrect job status counts.
427 for subdag_dir in {file.parent for file in path.glob("subdags/*/*.rescue*")}:
428 if failed_subdags is not None and subdag_dir.name not in failed_subdags: 428 ↛ 429line 428 didn't jump to line 429 because the condition on line 428 was never true
429 continue
430 subdag_dest = dest / subdag_dir.relative_to(path)
431 subdag_dest.mkdir(parents=True, exist_ok=False)
432 htc_backup_files_single_path(subdag_dir, subdag_dest)
435def htc_backup_files_single_path(src: str | os.PathLike, dest: str | os.PathLike) -> None:
436 """Move particular htc files to a different directory for later debugging.
438 Parameters
439 ----------
440 src : `str` or `os.PathLike`
441 Directory from which to back up particular files.
442 dest : `str` or `os.PathLike`
443 Directory to which particular files are moved.
445 Raises
446 ------
447 RuntimeError
448 If given dest directory matches given src directory.
449 OSError
450 If problems moving file.
451 FileNotFoundError
452 Item matching pattern in src directory isn't a file.
453 """
454 src = Path(src)
455 dest = Path(dest)
456 if dest.samefile(src):
457 raise RuntimeError(f"Destination directory is same as the source directory ({src})")
459 for patt in [
460 "*.info.*",
461 "*.dag.metrics",
462 "*.dag.nodes.log",
463 "*.node_status",
464 "wms_*.dag.post.out",
465 "wms_*.status.txt",
466 ]:
467 for source in src.glob(patt):
468 if source.is_file(): 468 ↛ 475line 468 didn't jump to line 475 because the condition on line 468 was always true
469 target = dest / source.relative_to(src)
470 try:
471 source.rename(target)
472 except OSError as exc:
473 raise type(exc)(f"Backing up '{source}' failed: {exc.strerror}") from None
474 else:
475 raise FileNotFoundError(f"Backing up '{source}' failed: not a file")
478def htc_escape(value):
479 """Escape characters in given value based upon HTCondor syntax.
481 Parameters
482 ----------
483 value : `~typing.Any`
484 Value that needs to have characters escaped if string.
486 Returns
487 -------
488 new_value : `~typing.Any`
489 Given value with characters escaped appropriate for HTCondor if string.
490 """
491 if isinstance(value, str):
492 newval = value.replace('"', '""').replace("'", "''").replace(""", '"')
493 else:
494 newval = value
496 return newval
499def htc_write_attribs(stream, attrs):
500 """Write job attributes in HTCondor format to writeable stream.
502 Parameters
503 ----------
504 stream : `~typing.TextIO`
505 Output text stream (typically an open file).
506 attrs : `dict`
507 HTCondor job attributes (dictionary of attribute key, value).
508 """
509 for key, value in attrs.items():
510 # Make sure strings are syntactically correct for HTCondor.
511 if isinstance(value, str): 511 ↛ 514line 511 didn't jump to line 514 because the condition on line 511 was always true
512 pval = f'"{htc_escape(value)}"'
513 else:
514 pval = value
516 print(f"+{key} = {pval}", file=stream)
519def htc_write_condor_file(
520 filename: str | os.PathLike, job_name: str, job: RestrictedDict, job_attrs: dict[str, Any]
521) -> None:
522 """Write an HTCondor submit file.
524 Parameters
525 ----------
526 filename : `str` or `os.PathLike`
527 Filename for the HTCondor submit file.
528 job_name : `str`
529 Job name to use in submit file.
530 job : `RestrictedDict`
531 Submit script information.
532 job_attrs : `dict`
533 Job attributes.
534 """
535 os.makedirs(os.path.dirname(filename), exist_ok=True)
536 with open(filename, "w") as fh:
537 for key, value in job.items():
538 if value is not None: 538 ↛ 537line 538 didn't jump to line 537 because the condition on line 538 was always true
539 if key in HTC_QUOTE_KEYS: # Assumes internal quotes are already escaped correctly
540 print(f'{key}="{value}"', file=fh)
541 else:
542 print(f"{key}={value}", file=fh)
543 for key in ["output", "error", "log"]:
544 if key not in job:
545 filename = f"{job_name}.$(Cluster).{'out' if key != 'log' else key}"
546 print(f"{key}={filename}", file=fh)
548 if job_attrs is not None:
549 htc_write_attribs(fh, job_attrs)
550 print("queue", file=fh)
553# To avoid doing the version check during every function call select
554# appropriate conversion function at the import time.
555#
556# Make sure that *each* version specific variant of the conversion function(s)
557# has the same signature after applying any changes!
558if HTC_VERSION < version.parse("8.9.8"): 558 ↛ 560line 558 didn't jump to line 560 because the condition on line 558 was never true
560 def htc_tune_schedd_args(**kwargs):
561 """Ensure that arguments for Schedd are version appropriate.
563 The old arguments: 'requirements' and 'attr_list' of
564 'Schedd.history()', 'Schedd.query()', and 'Schedd.xquery()' were
565 deprecated in favor of 'constraint' and 'projection', respectively,
566 starting from version 8.9.8. The function will convert "new" keyword
567 arguments to "old" ones.
569 Parameters
570 ----------
571 **kwargs
572 Any keyword arguments that Schedd.history(), Schedd.query(), and
573 Schedd.xquery() accepts.
575 Returns
576 -------
577 kwargs : `dict` [`str`, `~typing.Any`]
578 Keywords arguments that are guaranteed to work with the Python
579 HTCondor API.
581 Notes
582 -----
583 Function doesn't validate provided keyword arguments beyond converting
584 selected arguments to their version specific form. For example,
585 it won't remove keywords that are not supported by the methods
586 mentioned earlier.
587 """
588 translation_table = {
589 "constraint": "requirements",
590 "projection": "attr_list",
591 }
592 for new, old in translation_table.items():
593 try:
594 kwargs[old] = kwargs.pop(new)
595 except KeyError:
596 pass
597 return kwargs
599else:
601 def htc_tune_schedd_args(**kwargs):
602 """Ensure that arguments for Schedd are version appropriate.
604 This is the fallback function if no version specific alteration are
605 necessary. Effectively, a no-op.
607 Parameters
608 ----------
609 **kwargs
610 Any keyword arguments that Schedd.history(), Schedd.query(), and
611 Schedd.xquery() accepts.
613 Returns
614 -------
615 kwargs : `dict` [`str`, `~typing.Any`]
616 Keywords arguments that were passed to the function.
617 """
618 return kwargs
621def htc_query_history(schedds, **kwargs):
622 """Fetch history records from the condor_schedd daemon.
624 Parameters
625 ----------
626 schedds : `htcondor.Schedd`
627 HTCondor schedulers which to query for job information.
628 **kwargs
629 Any keyword arguments that Schedd.history() accepts.
631 Yields
632 ------
633 schedd_name : `str`
634 Name of the HTCondor scheduler managing the job queue.
635 job_ad : `dict` [`str`, `~typing.Any`]
636 A dictionary representing HTCondor ClassAd describing a job. It maps
637 job attributes names to values of the ClassAd expressions they
638 represent.
639 """
640 # If not set, provide defaults for positional arguments.
641 kwargs.setdefault("constraint", None)
642 kwargs.setdefault("projection", [])
643 kwargs = htc_tune_schedd_args(**kwargs)
644 for schedd_name, schedd in schedds.items():
645 for job_ad in schedd.history(**kwargs):
646 yield schedd_name, dict(job_ad)
649def htc_query_present(schedds, **kwargs):
650 """Query the condor_schedd daemon for job ads.
652 Parameters
653 ----------
654 schedds : `htcondor.Schedd`
655 HTCondor schedulers which to query for job information.
656 **kwargs
657 Any keyword arguments that Schedd.xquery() accepts.
659 Yields
660 ------
661 schedd_name : `str`
662 Name of the HTCondor scheduler managing the job queue.
663 job_ad : `dict` [`str`, `~typing.Any`]
664 A dictionary representing HTCondor ClassAd describing a job. It maps
665 job attributes names to values of the ClassAd expressions they
666 represent.
667 """
668 kwargs = htc_tune_schedd_args(**kwargs)
669 for schedd_name, schedd in schedds.items():
670 for job_ad in schedd.query(**kwargs):
671 yield schedd_name, dict(job_ad)
674def htc_version():
675 """Return the version given by the HTCondor API.
677 Returns
678 -------
679 version : `str`
680 HTCondor version as easily comparable string.
681 """
682 return str(HTC_VERSION)
685def htc_submit_dag(sub):
686 """Submit job for execution.
688 Parameters
689 ----------
690 sub : `htcondor.Submit`
691 An object representing a job submit description.
693 Returns
694 -------
695 schedd_job_info : `dict` [`str`, `dict` [`str`, \
696 `dict` [`str`, `~typing.Any`]]]
697 Information about jobs satisfying the search criteria where for each
698 Scheduler, local HTCondor job ids are mapped to their respective
699 classads.
700 """
701 coll = htcondor.Collector()
702 schedd_ad = coll.locate(htcondor.DaemonTypes.Schedd)
703 schedd = htcondor.Schedd(schedd_ad)
705 # If Schedd.submit() fails, the method will raise an exception. Usually,
706 # that implies issues with the HTCondor pool which BPS can't address.
707 # Hence, no effort is made to handle the exception.
708 submit_result = schedd.submit(sub)
710 # Sadly, the ClassAd from Schedd.submit() (see above) does not have
711 # 'GlobalJobId' so we need to run a regular query to get it anyway.
712 schedd_name = schedd_ad["Name"]
713 schedd_dag_info = condor_q(
714 constraint=f"ClusterId == {submit_result.cluster()}", schedds={schedd_name: schedd}
715 )
716 return schedd_dag_info
719def htc_create_submit_from_dag(
720 dag_filename: str, submit_options: dict[str, Any], dagman_conf_filename: str | os.PathLike | None = None
721) -> htcondor.Submit:
722 """Create a DAGMan job submit description.
724 Parameters
725 ----------
726 dag_filename : `str`
727 Name of file containing HTCondor DAG commands.
728 submit_options : `dict` [`str`, `~typing.Any`], optional
729 Contains extra options for command line (Value of None means flag).
730 dagman_conf_filename : `str` or `os.PathLike`, optional
731 Location of DAGMan configuration file, if Any. Defaults to no
732 file (None).
734 Returns
735 -------
736 sub : `htcondor.Submit`
737 An object representing a job submit description.
739 Notes
740 -----
741 Use with HTCondor versions which support htcondor.Submit.from_dag(),
742 i.e., 8.9.3 or newer.
743 """
744 # Config and environment variables do not seem to override -MaxIdle
745 # on the .dag.condor.sub's command line (broken in some 24.0.x versions).
746 # Explicitly forward them as a submit_option if either exists.
747 # Note: auto generated subdag submit files are still the -MaxIdle=1000
748 # in the broken versions.
749 if "MaxIdle" not in submit_options:
750 max_jobs_idle: int | None = None
751 config_var_name = "DAGMAN_MAX_JOBS_IDLE"
753 if dagman_conf_filename:
754 _LOG.debug("Checking DAGMan config file = %s", dagman_conf_filename)
755 with open(dagman_conf_filename) as fh:
756 for line in fh:
757 _LOG.debug("DAGMan config file line = %s", line)
758 parts = line.split("=")
759 _LOG.debug("DAGMan config file line parts = %s", parts)
760 if len(parts) == 2 and parts[0].strip() == config_var_name:
761 max_jobs_idle = int(parts[1].strip())
762 _LOG.debug("Found %s = %s", config_var_name, max_jobs_idle)
763 break
764 if max_jobs_idle is None:
765 if f"_CONDOR_{config_var_name}" in os.environ:
766 max_jobs_idle = int(os.environ[f"_CONDOR_{config_var_name}"])
767 elif config_var_name in htcondor.param:
768 max_jobs_idle = htcondor.param[config_var_name]
770 if max_jobs_idle:
771 submit_options["MaxIdle"] = max_jobs_idle
772 else:
773 _LOG.debug("MaxIdle already in submit_options: %s", submit_options)
775 _LOG.debug("Using submit_options = %s", submit_options)
776 return htcondor.Submit.from_dag(dag_filename, submit_options)
779def htc_create_submit_from_cmd(dag_filename, submit_options=None):
780 """Create a DAGMan job submit description.
782 Create a DAGMan job submit description by calling ``condor_submit_dag``
783 on given DAG description file.
785 Parameters
786 ----------
787 dag_filename : `str`
788 Name of file containing HTCondor DAG commands.
789 submit_options : `dict` [`str`, `~typing.Any`], optional
790 Contains extra options for command line (Value of None means flag).
792 Returns
793 -------
794 sub : `htcondor.Submit`
795 An object representing a job submit description.
797 Notes
798 -----
799 Use with HTCondor versions which do not support htcondor.Submit.from_dag(),
800 i.e., older than 8.9.3.
801 """
802 # Run command line condor_submit_dag command.
803 cmd = "condor_submit_dag -f -no_submit -notification never -autorescue 1 -UseDagDir -no_recurse "
805 if submit_options is not None:
806 for opt, val in submit_options.items():
807 cmd += f" -{opt} {val or ''}"
808 cmd += f"{dag_filename}"
810 process = subprocess.Popen(
811 cmd.split(), shell=False, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, encoding="utf-8"
812 )
813 process.wait()
815 if process.returncode != 0:
816 print(f"Exit code: {process.returncode}")
817 print(process.communicate()[0])
818 raise RuntimeError("Problems running condor_submit_dag")
820 return htc_create_submit_from_file(f"{dag_filename}.condor.sub")
823def htc_create_submit_from_file(submit_file):
824 """Parse a submission file.
826 Parameters
827 ----------
828 submit_file : `str`
829 Name of the HTCondor submit file.
831 Returns
832 -------
833 sub : `htcondor.Submit`
834 An object representing a job submit description.
835 """
836 descriptors = {}
837 with open(submit_file) as fh:
838 for line in fh:
839 line = line.strip()
840 if not line.startswith("#") and not line == "queue":
841 (key, val) = re.split(r"\s*=\s*", line, maxsplit=1)
842 descriptors[key] = val
844 # Avoid UserWarning: the line 'copy_to_spool = False' was
845 # unused by Submit object. Is it a typo?
846 try:
847 del descriptors["copy_to_spool"]
848 except KeyError:
849 pass
851 return htcondor.Submit(descriptors)
854def _htc_write_job_commands(stream, name, commands, node_type="JOB"):
855 """Output the DAGMan job lines for single job in DAG.
857 Parameters
858 ----------
859 stream : `~typing.TextIO`
860 Writeable text stream (typically an opened file).
861 name : `str`
862 Job name.
863 commands : `RestrictedDict`
864 DAG commands for a job.
865 node_type : `str`, optional
866 Type of DAGMan node (JOB, FINAL, SERVICE). Defaults to "JOB".
867 """
868 # Note: optional pieces of commands include a space at the beginning.
869 # also making sure values aren't empty strings as placeholders.
870 if "pre" in commands and commands["pre"]:
871 defer = ""
872 if "defer" in commands["pre"] and commands["pre"]["defer"]:
873 defer = f" DEFER {commands['pre']['defer']['status']} {commands['pre']['defer']['time']}"
875 debug = ""
876 if "debug" in commands["pre"] and commands["pre"]["debug"]:
877 debug = f" DEBUG {commands['pre']['debug']['filename']} {commands['pre']['debug']['type']}"
879 arguments = ""
880 if "arguments" in commands["pre"] and commands["pre"]["arguments"]:
881 arguments = f" {commands['pre']['arguments']}"
883 executable = commands["pre"]["executable"]
884 print(f"SCRIPT{defer}{debug} PRE {name} {executable}{arguments}", file=stream)
886 if "post" in commands and commands["post"]:
887 defer = ""
888 if "defer" in commands["post"] and commands["post"]["defer"]:
889 defer = f" DEFER {commands['post']['defer']['status']} {commands['post']['defer']['time']}"
891 debug = ""
892 if "debug" in commands["post"] and commands["post"]["debug"]:
893 debug = f" DEBUG {commands['post']['debug']['filename']} {commands['post']['debug']['type']}"
895 arguments = ""
896 if "arguments" in commands["post"] and commands["post"]["arguments"]:
897 arguments = f" {commands['post']['arguments']}"
899 executable = commands["post"]["executable"]
900 print(f"SCRIPT{defer}{debug} POST {name} {executable}{arguments}", file=stream)
902 if "vars" in commands and commands["vars"]:
903 for key, value in commands["vars"].items():
904 print(f'VARS {name} {key}="{htc_escape(value)}"', file=stream)
906 if "pre_skip" in commands and commands["pre_skip"]:
907 print(f"PRE_SKIP {name} {commands['pre_skip']}", file=stream)
909 # FINAL node cannot have a DAGMan retry, abort-dag-on, priority, category
910 if node_type != "FINAL":
911 if "retry" in commands and commands["retry"]:
912 print(f"RETRY {name} {commands['retry']}", end="", file=stream)
913 if "retry_unless_exit" in commands and commands["retry_unless_exit"]:
914 print(f" UNLESS-EXIT {commands['retry_unless_exit']}", end="", file=stream)
915 print("", file=stream) # Since previous prints don't include new line
917 if "abort_dag_on" in commands and commands["abort_dag_on"]:
918 print(
919 f"ABORT-DAG-ON {name} {commands['abort_dag_on']['node_exit']}"
920 f" RETURN {commands['abort_dag_on']['abort_exit']}",
921 file=stream,
922 )
924 if "priority" in commands and commands["priority"]:
925 print(
926 f"PRIORITY {name} {commands['priority']}",
927 file=stream,
928 )
931class HTCJob:
932 """HTCondor job for use in building DAG.
934 Parameters
935 ----------
936 name : `str`
937 Name of the job.
938 label : `str`
939 Label that can used for grouping or lookup.
940 initcmds : `RestrictedDict`
941 Initial job commands for submit file.
942 initdagcmds : `RestrictedDict`
943 Initial commands for job inside DAG.
944 initattrs : `dict`
945 Initial dictionary of job attributes.
946 """
948 def __init__(self, name, label=None, initcmds=(), initdagcmds=(), initattrs=None):
949 self.name = name
950 self.label = label
951 self.cmds = RestrictedDict(HTC_VALID_JOB_KEYS, initcmds)
952 self.dagcmds = RestrictedDict(HTC_VALID_JOB_DAG_KEYS, initdagcmds)
953 self.attrs = initattrs
954 self.subfile = None
955 self.subdir = None
956 self.subdag = None
958 def __str__(self):
959 return self.name
961 def add_job_cmds(self, new_commands):
962 """Add commands to Job (overwrite existing).
964 Parameters
965 ----------
966 new_commands : `dict`
967 Submit file commands to be added to Job.
968 """
969 self.cmds.update(new_commands)
971 def add_dag_cmds(self, new_commands):
972 """Add DAG commands to Job (overwrite existing).
974 Parameters
975 ----------
976 new_commands : `dict`
977 DAG file commands to be added to Job.
978 """
979 self.dagcmds.update(new_commands)
981 def add_job_attrs(self, new_attrs):
982 """Add attributes to Job (overwrite existing).
984 Parameters
985 ----------
986 new_attrs : `dict`
987 Attributes to be added to Job.
988 """
989 if self.attrs is None:
990 self.attrs = {}
991 if new_attrs:
992 self.attrs.update(new_attrs)
994 def write_submit_file(self, submit_path: str | os.PathLike) -> None:
995 """Write job description to submit file.
997 Parameters
998 ----------
999 submit_path : `str` or `os.PathLike`
1000 Prefix path for the submit file.
1001 """
1002 if not self.subfile:
1003 self.subfile = f"{self.name}.sub"
1005 subfile = self.subfile
1006 if self.subdir:
1007 subfile = Path(self.subdir) / subfile
1009 subfile = Path(os.path.expandvars(subfile))
1010 if not subfile.is_absolute():
1011 subfile = Path(submit_path) / subfile
1012 if not subfile.exists():
1013 _LOG.debug("Writing subfile: %s", subfile)
1014 htc_write_condor_file(subfile, self.name, self.cmds, self.attrs)
1015 else:
1016 _LOG.debug("Using existing subfile: %s", subfile)
1018 def write_dag_commands(self, stream, dag_rel_path, command_name="JOB"):
1019 """Write DAG commands for single job to output stream.
1021 Parameters
1022 ----------
1023 stream : `~typing.TextIO`
1024 Output Stream.
1025 dag_rel_path : `str`
1026 Relative path of dag to submit directory.
1027 command_name : `str`
1028 Name of the DAG command (e.g., JOB, FINAL).
1029 """
1030 subfile = os.path.expandvars(self.subfile)
1032 # JOB NodeName SubmitDescription [DIR directory] [NOOP] [DONE]
1033 job_line = f'{command_name} {self.name} "{subfile}"'
1034 if "dir" in self.dagcmds:
1035 dir_val = self.dagcmds["dir"]
1036 if dag_rel_path:
1037 dir_val = os.path.join(dag_rel_path, dir_val)
1038 job_line += f' DIR "{dir_val}"'
1039 if self.dagcmds.get("noop", False):
1040 job_line += " NOOP"
1042 print(job_line, file=stream)
1043 if self.dagcmds:
1044 _htc_write_job_commands(stream, self.name, self.dagcmds, command_name)
1046 def dump(self, fh):
1047 """Dump job information to output stream.
1049 Parameters
1050 ----------
1051 fh : `~typing.TextIO`
1052 Output stream.
1053 """
1054 printer = pprint.PrettyPrinter(indent=4, stream=fh)
1055 printer.pprint(self.name)
1056 printer.pprint(self.cmds)
1057 printer.pprint(self.attrs)
1060class HTCDag(networkx.DiGraph):
1061 """HTCondor DAG.
1063 Parameters
1064 ----------
1065 data : `~typing.Any`
1066 Initial graph data of any format that is supported
1067 by the to_network_graph() function.
1068 name : `str`
1069 Name for DAG.
1070 """
1072 def __init__(self, data=None, name=""):
1073 super().__init__(data=data, name=name)
1075 self.graph["attr"] = {}
1076 self.graph["run_id"] = None
1077 self.graph["submit_path"] = None
1078 self.graph["final_job"] = None
1079 self.graph["service_job"] = None
1080 self.graph["submit_options"] = {}
1082 def __str__(self):
1083 """Represent basic DAG info as string.
1085 Returns
1086 -------
1087 info : `str`
1088 String containing basic DAG info.
1089 """
1090 return f"{self.graph['name']} {len(self)}"
1092 def add_attribs(self, attribs=None):
1093 """Add attributes to the DAG.
1095 Parameters
1096 ----------
1097 attribs : `dict`
1098 DAG attributes.
1099 """
1100 if attribs is not None: 1100 ↛ exitline 1100 didn't return from function 'add_attribs' because the condition on line 1100 was always true
1101 self.graph["attr"].update(attribs)
1103 def add_job(self, job, parent_names=None, child_names=None):
1104 """Add an HTCJob to the HTCDag.
1106 Parameters
1107 ----------
1108 job : `HTCJob`
1109 HTCJob to add to the HTCDag.
1110 parent_names : `~collections.abc.Iterable` [`str`], optional
1111 Names of parent jobs.
1112 child_names : `~collections.abc.Iterable` [`str`], optional
1113 Names of child jobs.
1114 """
1115 assert isinstance(job, HTCJob)
1116 _LOG.debug("Adding job %s to dag", job.name)
1118 # Add dag level attributes to each job
1119 job.add_job_attrs(self.graph["attr"])
1121 self.add_node(job.name, data=job)
1123 if parent_names is not None: 1123 ↛ 1124line 1123 didn't jump to line 1124 because the condition on line 1123 was never true
1124 self.add_job_relationships(parent_names, [job.name])
1126 if child_names is not None: 1126 ↛ 1127line 1126 didn't jump to line 1127 because the condition on line 1126 was never true
1127 self.add_job_relationships(child_names, [job.name])
1129 def add_job_relationships(self, parents, children):
1130 """Add DAG edge between parents and children jobs.
1132 Parameters
1133 ----------
1134 parents : `list` [`str`]
1135 Contains parent job name(s).
1136 children : `list` [`str`]
1137 Contains children job name(s).
1138 """
1139 self.add_edges_from(itertools.product(parents, children))
1141 def add_final_job(self, job):
1142 """Add an HTCJob for the FINAL job in HTCDag.
1144 Parameters
1145 ----------
1146 job : `HTCJob`
1147 HTCJob to add to the HTCDag as a FINAL job.
1148 """
1149 # Add dag level attributes to each job
1150 job.add_job_attrs(self.graph["attr"])
1152 self.graph["final_job"] = job
1154 def add_service_job(self, job):
1155 """Add an HTCJob for the SERVICE job in HTCDag.
1157 Parameters
1158 ----------
1159 job : `HTCJob`
1160 HTCJob to add to the HTCDag as a SERVICE job.
1161 """
1162 # Add dag level attributes to each job
1163 job.add_job_attrs(self.graph["attr"])
1165 self.graph["service_job"] = job
1167 def del_job(self, job_name):
1168 """Delete the job from the DAG.
1170 Parameters
1171 ----------
1172 job_name : `str`
1173 Name of job in DAG to delete.
1174 """
1175 # Reconnect edges around node to delete
1176 parents = self.predecessors(job_name)
1177 children = self.successors(job_name)
1178 self.add_edges_from(itertools.product(parents, children))
1180 # Delete job node (which deletes its edges).
1181 self.remove_node(job_name)
1183 def write(self, submit_path, job_subdir="", dag_subdir="", dag_rel_path=""):
1184 """Write DAG to a file.
1186 Parameters
1187 ----------
1188 submit_path : `str`
1189 Prefix path for all outputs.
1190 job_subdir : `str`, optional
1191 Template for job subdir (submit_path + job_subdir).
1192 dag_subdir : `str`, optional
1193 DAG subdir (submit_path + dag_subdir).
1194 dag_rel_path : `str`, optional
1195 Prefix to job_subdir for jobs inside subdag.
1196 """
1197 self.graph["submit_path"] = submit_path
1198 self.graph["dag_filename"] = os.path.join(dag_subdir, f"{self.graph['name']}.dag")
1199 full_filename = os.path.join(submit_path, self.graph["dag_filename"])
1200 os.makedirs(os.path.dirname(full_filename), exist_ok=True)
1202 try:
1203 dagman_config_path = Path(self.graph["attr"]["bps_wms_config_path"])
1204 except KeyError:
1205 dagman_config_path = None
1206 with open(full_filename, "w") as fh:
1207 if dagman_config_path is not None:
1208 fh.write(f"CONFIG {dag_rel_path / dagman_config_path}\n")
1210 for name, nodeval in self.nodes().items():
1211 try:
1212 job = nodeval["data"]
1213 except KeyError:
1214 _LOG.error("Job %s doesn't have data (keys: %s).", name, nodeval.keys())
1215 raise
1216 if job.subdag: 1216 ↛ 1217line 1216 didn't jump to line 1217 because the condition on line 1216 was never true
1217 if job.subfile:
1218 this_dag_rel_path = ""
1219 else:
1220 this_dag_rel_path = "../.."
1221 dag_subdir = f"subdags/{job.name}"
1222 if "dir" in job.dagcmds:
1223 subdir = job.dagcmds["dir"]
1224 else:
1225 subdir = job_subdir
1226 if dagman_config_path is not None:
1227 job.subdag.add_attribs({"bps_wms_config_path": str(dagman_config_path)})
1228 job.subdag.write(submit_path, subdir, dag_subdir, this_dag_rel_path)
1229 fh.write(f"SUBDAG EXTERNAL {job.name} {Path(job.subdag.graph['dag_filename']).name}")
1230 if dag_subdir:
1231 fh.write(f" DIR {dag_subdir}")
1232 fh.write("\n")
1233 if job.dagcmds:
1234 _htc_write_job_commands(fh, job.name, job.dagcmds)
1235 else:
1236 job.write_submit_file(submit_path)
1237 job.write_dag_commands(fh, dag_rel_path)
1239 for edge in self.edges():
1240 print(f"PARENT {edge[0]} CHILD {edge[1]}", file=fh)
1242 if self.graph.get("write_dot", False):
1243 print(f"DOT {self.name}.dot", file=fh)
1245 print(f"NODE_STATUS_FILE {self.name}.node_status", file=fh)
1247 # Add bps attributes to dag submission
1248 for key, value in self.graph["attr"].items():
1249 print(f'SET_JOB_ATTR {key}= "{htc_escape(value)}"', file=fh)
1251 # Add special nodes if any.
1252 special_jobs = {
1253 "FINAL": self.graph["final_job"],
1254 "SERVICE": self.graph["service_job"],
1255 }
1256 for dagcmd, job in special_jobs.items():
1257 if job is not None:
1258 job.write_submit_file(submit_path)
1259 job.write_dag_commands(fh, dag_rel_path, dagcmd)
1261 def dump(self, fh):
1262 """Dump DAG info to output stream.
1264 Parameters
1265 ----------
1266 fh : `typing.IO`
1267 Where to dump DAG info as text.
1268 """
1269 for key, value in self.graph:
1270 print(f"{key}={value}", file=fh)
1271 for name, data in self.nodes().items():
1272 print(f"{name}:", file=fh)
1273 data.dump(fh)
1274 for edge in self.edges():
1275 print(f"PARENT {edge[0]} CHILD {edge[1]}", file=fh)
1276 if self.graph["final_job"]:
1277 print(f"FINAL {self.graph['final_job'].name}:", file=fh)
1278 self.graph["final_job"].dump(fh)
1280 def write_dot(self, filename):
1281 """Write a dot version of the DAG.
1283 Parameters
1284 ----------
1285 filename : `str`
1286 Name of the dot file.
1287 """
1288 pos = networkx.nx_agraph.graphviz_layout(self)
1289 networkx.draw(self, pos=pos)
1290 networkx.drawing.nx_pydot.write_dot(self, filename)
1293def condor_q(constraint=None, schedds=None, **kwargs):
1294 """Get information about the jobs in the HTCondor job queue(s).
1296 Parameters
1297 ----------
1298 constraint : `str`, optional
1299 Constraints to be passed to job query.
1300 schedds : `dict` [`str`, `htcondor.Schedd`], optional
1301 HTCondor schedulers which to query for job information. If None
1302 (default), the query will be run against local scheduler only.
1303 **kwargs : `~typing.Any`
1304 Additional keyword arguments that need to be passed to the internal
1305 query method.
1307 Returns
1308 -------
1309 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str`, `~typing.Any`]]]
1310 Information about jobs satisfying the search criteria where for each
1311 Scheduler, local HTCondor job ids are mapped to their respective
1312 classads.
1313 """
1314 return condor_query(constraint, schedds, htc_query_present, **kwargs)
1317def condor_history(constraint=None, schedds=None, **kwargs):
1318 """Get information about the jobs from HTCondor history records.
1320 Parameters
1321 ----------
1322 constraint : `str`, optional
1323 Constraints to be passed to job query.
1324 schedds : `dict` [`str`, `htcondor.Schedd`], optional
1325 HTCondor schedulers which to query for job information. If None
1326 (default), the query will be run against the history file of
1327 the local scheduler only.
1328 **kwargs : `~typing.Any`
1329 Additional keyword arguments that need to be passed to the internal
1330 query method.
1332 Returns
1333 -------
1334 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str`, `~typing.Any`]]]
1335 Information about jobs satisfying the search criteria where for each
1336 Scheduler, local HTCondor job ids are mapped to their respective
1337 classads.
1338 """
1339 return condor_query(constraint, schedds, htc_query_history, **kwargs)
1342def condor_query(constraint=None, schedds=None, query_func=htc_query_present, **kwargs):
1343 """Get information about HTCondor jobs.
1345 Parameters
1346 ----------
1347 constraint : `str`, optional
1348 Constraints to be passed to job query.
1349 schedds : `dict` [`str`, `htcondor.Schedd`], optional
1350 HTCondor schedulers which to query for job information. If None
1351 (default), the query will be run against the history file of
1352 the local scheduler only.
1353 query_func : `~collections.abc.Callable`
1354 An query function which takes following arguments:
1356 - ``schedds``: Schedulers to query (`list` [`htcondor.Schedd`]).
1357 - ``**kwargs``: Keyword arguments that will be passed to the query
1358 function.
1359 **kwargs : `~typing.Any`
1360 Additional keyword arguments that need to be passed to the query
1361 method.
1363 Returns
1364 -------
1365 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str`, `~typing.Any`]]]
1366 Information about jobs satisfying the search criteria where for each
1367 Scheduler, local HTCondor job ids are mapped to their respective
1368 classads.
1369 """
1370 if not schedds:
1371 coll = htcondor.Collector()
1372 schedd_ad = coll.locate(htcondor.DaemonTypes.Schedd)
1373 schedds = {schedd_ad["Name"]: htcondor.Schedd(schedd_ad)}
1375 # Make sure that 'ClusterId' and 'ProcId' attributes are always included
1376 # in the job classad. They are needed to construct the job id.
1377 added_attrs = set()
1378 if "projection" in kwargs and kwargs["projection"]:
1379 requested_attrs = set(kwargs["projection"])
1380 required_attrs = {"ClusterId", "ProcId"}
1381 added_attrs = required_attrs - requested_attrs
1382 for attr in added_attrs:
1383 kwargs["projection"].append(attr)
1385 unwanted_attrs = {"Env", "Environment"} | added_attrs
1386 job_info = defaultdict(dict)
1387 for schedd_name, job_ad in query_func(schedds, constraint=constraint, **kwargs):
1388 id_ = f"{job_ad['ClusterId']}.{job_ad['ProcId']}"
1389 for attr in set(job_ad) & unwanted_attrs:
1390 del job_ad[attr]
1391 job_info[schedd_name][id_] = job_ad
1392 _LOG.debug("query returned %d jobs", sum(len(val) for val in job_info.values()))
1394 # Restore the list of the requested attributes to its original value
1395 # if needed.
1396 if added_attrs:
1397 for attr in added_attrs:
1398 kwargs["projection"].remove(attr)
1400 # When returning the results filter out entries for schedulers with no jobs
1401 # matching the search criteria.
1402 return {key: val for key, val in job_info.items() if val}
1405def condor_search(constraint=None, hist=None, schedds=None):
1406 """Search for running and finished jobs satisfying given criteria.
1408 Parameters
1409 ----------
1410 constraint : `str`, optional
1411 Constraints to be passed to job query.
1412 hist : `float`
1413 Limit history search to this many days.
1414 schedds : `dict` [`str`, `htcondor.Schedd`], optional
1415 The list of the HTCondor schedulers which to query for job information.
1416 If None (default), only the local scheduler will be queried.
1418 Returns
1419 -------
1420 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str` `~typing.Any`]]]
1421 Information about jobs satisfying the search criteria where for each
1422 Scheduler, local HTCondor job ids are mapped to their respective
1423 classads.
1424 """
1425 if not schedds:
1426 coll = htcondor.Collector()
1427 schedd_ad = coll.locate(htcondor.DaemonTypes.Schedd)
1428 schedds = {schedd_ad["Name"]: htcondor.Schedd(locate_ad=schedd_ad)}
1430 job_info = condor_q(constraint=constraint, schedds=schedds)
1431 if hist is not None:
1432 _LOG.debug("Searching history going back %s days", hist)
1433 epoch = (datetime.now() - timedelta(days=hist)).timestamp()
1434 constraint += f" && (CompletionDate >= {epoch} || JobFinishedHookDone >= {epoch})"
1435 hist_info = condor_history(constraint, schedds=schedds)
1436 update_job_info(job_info, hist_info)
1437 return job_info
1440def condor_status(constraint=None, coll=None):
1441 """Get information about HTCondor pool.
1443 Parameters
1444 ----------
1445 constraint : `str`, optional
1446 Constraints to be passed to the query.
1447 coll : `htcondor.Collector`, optional
1448 Object representing HTCondor collector daemon.
1450 Returns
1451 -------
1452 pool_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1453 Mapping between HTCondor slot names and slot information (classAds).
1454 """
1455 if coll is None:
1456 coll = htcondor.Collector()
1457 try:
1458 pool_ads = coll.query(constraint=constraint)
1459 except OSError as ex:
1460 raise RuntimeError(f"Problem querying the Collector. (Constraint='{constraint}')") from ex
1462 pool_info = {}
1463 for slot in pool_ads:
1464 pool_info[slot["name"]] = dict(slot)
1465 _LOG.debug("condor_status returned %d ads", len(pool_info))
1466 return pool_info
1469def update_job_info(job_info, other_info):
1470 """Update results of a job query with results from another query.
1472 Parameters
1473 ----------
1474 job_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1475 Results of the job query that needs to be updated.
1476 other_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1477 Results of the other job query.
1479 Returns
1480 -------
1481 job_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1482 The updated results.
1483 """
1484 for schedd_name, others in other_info.items():
1485 try:
1486 jobs = job_info[schedd_name]
1487 except KeyError:
1488 job_info[schedd_name] = others
1489 else:
1490 for id_, ad in others.items():
1491 jobs.setdefault(id_, {}).update(ad)
1492 return job_info
1495def count_jobs_in_single_dag(
1496 filename: str | os.PathLike,
1497) -> tuple[Counter[str], dict[str, str], dict[str, WmsNodeType]]:
1498 """Build bps_run_summary string from dag file.
1500 Parameters
1501 ----------
1502 filename : `str`
1503 Path that includes dag file for a run.
1505 Returns
1506 -------
1507 counts : `Counter` [`str`]
1508 Semi-colon separated list of job labels and counts.
1509 (Same format as saved in dag classad).
1510 job_name_to_label : `dict` [`str`, `str`]
1511 Mapping of job names to job labels.
1512 job_name_to_type : `dict` [`str`, `lsst.ctrl.bps.htcondor.WmsNodeType`]
1513 Mapping of job names to job types
1514 (e.g., payload, final, service).
1515 """
1516 # Later code depends upon insertion order
1517 counts: Counter = Counter() # counts of payload jobs per label
1518 job_name_to_label: dict[str, str] = {}
1519 job_name_to_type: dict[str, WmsNodeType] = {}
1520 with open(filename) as fh:
1521 for line in fh:
1522 # Skip any line that contains commands irrelevant to job counting.
1523 if not line.startswith(
1524 (
1525 "JOB",
1526 "FINAL",
1527 "SERVICE",
1528 "SUBDAG EXTERNAL",
1529 )
1530 ):
1531 continue
1533 m = re.match(
1534 r"(?P<command>JOB|FINAL|SERVICE|SUBDAG EXTERNAL)\s+"
1535 r'(?P<jobname>(?P<wms>wms_)?\S+)\s+"?(?P<subfile>\S+)"?\s*'
1536 r'(DIR "?(?P<dir>[^\s"]+)"?)?\s*(?P<noop>NOOP)?',
1537 line,
1538 )
1539 if m: 1539 ↛ 1592line 1539 didn't jump to line 1592 because the condition on line 1539 was always true
1540 job_name = m.group("jobname")
1541 job_type = WmsNodeType.UNKNOWN
1542 name_parts = job_name.split("_")
1544 label = ""
1545 if m.group("dir"):
1546 dir_match = re.search(r"jobs/([^\s/]+)", m.group("dir"))
1547 if dir_match:
1548 label = dir_match.group(1)
1549 else:
1550 _LOG.debug("Parse DAG: unparsed dir = %s", line)
1551 elif m.group("subfile"): 1551 ↛ 1558line 1551 didn't jump to line 1558 because the condition on line 1551 was always true
1552 subfile_match = re.search(r"jobs/([^\s/]+)", m.group("subfile"))
1553 if subfile_match:
1554 label = m.group("subfile").split("/")[1]
1555 else:
1556 label = pegasus_name_to_label(job_name)
1558 match m.group("command"):
1559 case "JOB":
1560 if m.group("noop"):
1561 job_type = WmsNodeType.NOOP
1562 # wms_noop_label
1563 label = name_parts[2]
1564 elif m.group("wms"):
1565 if name_parts[1] == "check": 1565 ↛ 1570line 1565 didn't jump to line 1570 because the condition on line 1565 was always true
1566 job_type = WmsNodeType.SUBDAG_CHECK
1567 # wms_check_status_wms_group_label
1568 label = name_parts[5]
1569 else:
1570 _LOG.warning(
1571 "Unexpected skipping of dag line due to unknown wms job: %s", line
1572 )
1573 else:
1574 job_type = WmsNodeType.PAYLOAD
1575 if label == "init": 1575 ↛ 1576line 1575 didn't jump to line 1576 because the condition on line 1575 was never true
1576 label = "pipetaskInit"
1577 counts[label] += 1
1578 case "FINAL":
1579 job_type = WmsNodeType.FINAL
1580 counts[label] += 1 # final counts a payload job.
1581 case "SERVICE":
1582 job_type = WmsNodeType.SERVICE
1583 case "SUBDAG EXTERNAL": 1583 ↛ 1587line 1583 didn't jump to line 1587 because the pattern on line 1583 always matched
1584 job_type = WmsNodeType.SUBDAG
1585 label = name_parts[2]
1587 job_name_to_label[job_name] = label
1588 job_name_to_type[job_name] = job_type
1589 else:
1590 # The line should, but didn't match the pattern above. Probably
1591 # problems with regex.
1592 _LOG.warning("Unexpected skipping of dag line: %s", line)
1594 return counts, job_name_to_label, job_name_to_type
1597def summarize_dag(dir_name: str) -> tuple[str, dict[str, str], dict[str, WmsNodeType]]:
1598 """Build bps_run_summary string from dag file.
1600 Parameters
1601 ----------
1602 dir_name : `str`
1603 Path that includes dag file for a run.
1605 Returns
1606 -------
1607 summary : `str`
1608 Semi-colon separated list of job labels and counts
1609 (Same format as saved in dag classad).
1610 job_name_to_label : `dict` [`str`, `str`]
1611 Mapping of job names to job labels.
1612 job_name_to_type : `dict` [`str`, `lsst.ctrl.bps.htcondor.WmsNodeType`]
1613 Mapping of job names to job types
1614 (e.g., payload, final, service).
1615 """
1616 # Later code depends upon insertion order
1617 counts: Counter[str] = Counter() # counts of payload jobs per label
1618 job_name_to_label: dict[str, str] = {}
1619 job_name_to_type: dict[str, WmsNodeType] = {}
1620 for filename in Path(dir_name).glob("*.dag"):
1621 single_counts, single_job_name_to_label, single_job_name_to_type = count_jobs_in_single_dag(filename)
1622 counts += single_counts
1623 _update_dicts(job_name_to_label, single_job_name_to_label)
1624 _update_dicts(job_name_to_type, single_job_name_to_type)
1626 for filename in Path(dir_name).glob("subdags/*/*.dag"):
1627 single_counts, single_job_name_to_label, single_job_name_to_type = count_jobs_in_single_dag(filename)
1628 counts += single_counts
1629 _update_dicts(job_name_to_label, single_job_name_to_label)
1630 _update_dicts(job_name_to_type, single_job_name_to_type)
1632 summary = ";".join([f"{name}:{counts[name]}" for name in counts])
1633 _LOG.debug("summarize_dag: %s %s %s", summary, job_name_to_label, job_name_to_type)
1634 return summary, job_name_to_label, job_name_to_type
1637def pegasus_name_to_label(name):
1638 """Convert pegasus job name to a label for the report.
1640 Parameters
1641 ----------
1642 name : `str`
1643 Name of job.
1645 Returns
1646 -------
1647 label : `str`
1648 Label for job.
1649 """
1650 label = "UNK"
1651 if name.startswith("create_dir") or name.startswith("stage_in") or name.startswith("stage_out"): 1651 ↛ 1652line 1651 didn't jump to line 1652 because the condition on line 1651 was never true
1652 label = "pegasus"
1653 else:
1654 m = re.match(r"pipetask_(\d+_)?([^_]+)", name)
1655 if m: 1655 ↛ 1656line 1655 didn't jump to line 1656 because the condition on line 1655 was never true
1656 label = m.group(2)
1657 if label == "init":
1658 label = "pipetaskInit"
1660 return label
1663def read_single_dag_status(filename: str | os.PathLike) -> dict[str, Any]:
1664 """Read the node status file for DAG summary information.
1666 Parameters
1667 ----------
1668 filename : `str` or `Path.pathlib`
1669 Node status filename.
1671 Returns
1672 -------
1673 dag_ad : `dict` [`str`, `~typing.Any`]
1674 DAG summary information.
1675 """
1676 dag_ad: dict[str, Any] = {}
1678 # While this is probably more up to date than dag classad, only read from
1679 # file if need to.
1680 try:
1681 node_stat_file = Path(filename)
1682 _LOG.debug("Reading Node Status File %s", node_stat_file)
1683 with open(node_stat_file) as infh:
1684 dag_ad = dict(classad.parseNext(infh)) # pylint: disable=E1101
1686 if not dag_ad: 1686 ↛ 1688line 1686 didn't jump to line 1688 because the condition on line 1686 was never true
1687 # Pegasus check here
1688 metrics_file = node_stat_file.with_suffix(".dag.metrics")
1689 if metrics_file.exists():
1690 with open(metrics_file) as infh:
1691 metrics = json.load(infh)
1692 dag_ad["NodesTotal"] = metrics.get("jobs", 0)
1693 dag_ad["NodesFailed"] = metrics.get("jobs_failed", 0)
1694 dag_ad["NodesDone"] = metrics.get("jobs_succeeded", 0)
1695 metrics_file = node_stat_file.with_suffix(".metrics")
1696 with open(metrics_file) as infh:
1697 metrics = json.load(infh)
1698 dag_ad["NodesTotal"] = metrics["wf_metrics"]["total_jobs"]
1699 except (OSError, PermissionError):
1700 pass
1702 _LOG.debug("read_dag_status: %s", dag_ad)
1703 return dag_ad
1706def read_dag_status(wms_path: str | os.PathLike) -> dict[str, Any]:
1707 """Read the node status file for DAG summary information.
1709 Parameters
1710 ----------
1711 wms_path : `str` or `os.PathLike`
1712 Path that includes node status file for a run.
1714 Returns
1715 -------
1716 dag_ad : `dict` [`str`, `~typing.Any`]
1717 DAG summary information, counts summed across any subdags.
1718 """
1719 dag_ads: dict[str, Any] = {}
1720 path = Path(wms_path)
1721 try:
1722 node_stat_file = next(path.glob("*.node_status"))
1723 except StopIteration as exc:
1724 raise FileNotFoundError(f"DAGMan node status not found in {wms_path}") from exc
1726 dag_ads = read_single_dag_status(node_stat_file)
1728 for node_stat_file in path.glob("subdags/*/*.node_status"):
1729 dag_ad = read_single_dag_status(node_stat_file)
1730 dag_ads["JobProcsHeld"] += dag_ad.get("JobProcsHeld", 0)
1731 dag_ads["NodesPost"] += dag_ad.get("NodesPost", 0)
1732 dag_ads["JobProcsIdle"] += dag_ad.get("JobProcsIdle", 0)
1733 dag_ads["NodesTotal"] += dag_ad.get("NodesTotal", 0)
1734 dag_ads["NodesFailed"] += dag_ad.get("NodesFailed", 0)
1735 dag_ads["NodesDone"] += dag_ad.get("NodesDone", 0)
1736 dag_ads["NodesQueued"] += dag_ad.get("NodesQueued", 0)
1737 dag_ads["NodesPre"] += dag_ad.get("NodesReady", 0)
1738 dag_ads["NodesFutile"] += dag_ad.get("NodesFutile", 0)
1739 dag_ads["NodesUnready"] += dag_ad.get("NodesUnready", 0)
1741 return dag_ads
1744def read_single_node_status(filename: str | os.PathLike, init_fake_id: int) -> dict[str, Any]:
1745 """Read entire node status file.
1747 Parameters
1748 ----------
1749 filename : `str` or `pathlib.Path`
1750 Node status filename.
1751 init_fake_id : `int`
1752 Initial fake id value.
1754 Returns
1755 -------
1756 jobs : `dict` [`str`, `~typing.Any`]
1757 DAG summary information compiled from the node status file combined
1758 with the information found in the node event log.
1760 Currently, if the same job attribute is found in both files, its value
1761 from the event log takes precedence over the value from the node status
1762 file.
1763 """
1764 filename = Path(filename)
1766 # Get jobid info from other places to fill in gaps in info from node_status
1767 _, job_name_to_label, job_name_to_type = count_jobs_in_single_dag(filename.with_suffix(".dag"))
1768 loginfo: dict[str, dict[str, Any]] = {}
1769 wms_workflow_id = MISSING_ID
1770 try:
1771 wms_workflow_id, _ = read_single_dag_log(filename.with_suffix(".dag.dagman.log"))
1772 loginfo = read_single_dag_nodes_log(filename.with_suffix(".dag.nodes.log"))
1773 except (OSError, PermissionError):
1774 pass
1776 job_name_to_id: dict[str, str] = {}
1777 _LOG.debug("loginfo = %s", loginfo)
1778 log_job_name_to_id: dict[str, str] = {}
1779 for job_id, job_info in loginfo.items():
1780 if "LogNotes" in job_info: 1780 ↛ 1779line 1780 didn't jump to line 1779 because the condition on line 1780 was always true
1781 m = re.match(r"DAG Node: (\S+)", job_info["LogNotes"])
1782 if m: 1782 ↛ 1779line 1782 didn't jump to line 1779 because the condition on line 1782 was always true
1783 job_name = m.group(1)
1784 log_job_name_to_id[job_name] = job_id
1785 job_info["DAGNodeName"] = job_name
1786 job_info["wms_node_type"] = job_name_to_type[job_name]
1787 job_info["bps_job_label"] = job_name_to_label[job_name]
1789 jobs = {}
1790 fake_id = init_fake_id # For nodes that do not yet have a job id, give fake one
1791 try:
1792 with open(filename) as fh:
1793 for ad in classad.parseAds(fh):
1794 match ad["Type"]:
1795 case "DagStatus":
1796 # Skip DAG summary.
1797 pass
1798 case "NodeStatus":
1799 job_name = ad["Node"]
1800 if job_name in job_name_to_label: 1800 ↛ 1802line 1800 didn't jump to line 1802 because the condition on line 1800 was always true
1801 job_label = job_name_to_label[job_name]
1802 elif "_" in job_name:
1803 job_label = job_name.split("_")[1]
1804 else:
1805 job_label = job_name
1807 job = dict(ad)
1808 if job_name in log_job_name_to_id:
1809 job_id = str(log_job_name_to_id[job_name])
1810 _update_dicts(job, loginfo[job_id])
1811 else:
1812 job_id = str(fake_id)
1813 job = dict(ad)
1814 fake_id -= 1
1815 jobs[job_id] = job
1816 job_name_to_id[job_name] = job_id
1818 # Make job info as if came from condor_q.
1819 job["ClusterId"] = int(float(job_id))
1820 job["DAGManJobID"] = wms_workflow_id
1821 job["DAGNodeName"] = job_name
1822 job["bps_job_label"] = job_label
1823 job["wms_node_type"] = job_name_to_type[job_name]
1825 case "StatusEnd": 1825 ↛ 1828line 1825 didn't jump to line 1828 because the pattern on line 1825 always matched
1826 # Skip node status file "epilog".
1827 pass
1828 case _:
1829 _LOG.debug(
1830 "Ignoring unknown classad type '%s' in the node status file '%s'",
1831 ad["Type"],
1832 filename,
1833 )
1834 except (OSError, PermissionError):
1835 pass
1837 # Check for missing jobs (e.g., submission failure or not submitted yet)
1838 # Use dag info to create job placeholders
1839 for name in set(job_name_to_label) - set(job_name_to_id):
1840 if name in log_job_name_to_id: # job was in nodes.log, but not node_status
1841 job_id = str(log_job_name_to_id[name])
1842 job = dict(loginfo[job_id])
1843 else:
1844 job_id = str(fake_id)
1845 fake_id -= 1
1846 job = {}
1847 job["NodeStatus"] = NodeStatus.NOT_READY
1849 job["ClusterId"] = int(float(job_id))
1850 job["ProcId"] = 0
1851 job["DAGManJobID"] = wms_workflow_id
1852 job["DAGNodeName"] = name
1853 job["bps_job_label"] = job_name_to_label[name]
1854 job["wms_node_type"] = job_name_to_type[name]
1855 jobs[f"{job['ClusterId']}.{job['ProcId']}"] = job
1857 for job_info in jobs.values():
1858 job_info["from_dag_job"] = f"wms_{filename.stem}"
1860 return jobs
1863def read_node_status(wms_path: str | os.PathLike) -> dict[str, dict[str, Any]]:
1864 """Read entire node status file.
1866 Parameters
1867 ----------
1868 wms_path : `str` or `os.PathLike`
1869 Path that includes node status file for a run.
1871 Returns
1872 -------
1873 jobs : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1874 DAG summary information compiled from the node status file combined
1875 with the information found in the node event log.
1877 Currently, if the same job attribute is found in both files, its value
1878 from the event log takes precedence over the value from the node status
1879 file.
1880 """
1881 jobs: dict[str, dict[str, Any]] = {}
1882 init_fake_id = -1
1884 # subdags may not have run so wouldn't have node_status file
1885 # use dag files and let read_single_node_status handle missing
1886 # node_status file.
1887 for dag_filename in Path(wms_path).glob("*.dag"):
1888 filename = dag_filename.with_suffix(".node_status")
1889 info = read_single_node_status(filename, init_fake_id)
1890 init_fake_id -= len(info)
1891 _update_dicts(jobs, info)
1893 for dag_filename in Path(wms_path).glob("subdags/*/*.dag"):
1894 filename = dag_filename.with_suffix(".node_status")
1895 info = read_single_node_status(filename, init_fake_id)
1896 init_fake_id -= len(info)
1897 _update_dicts(jobs, info)
1899 # Propagate pruned from subdags to jobs
1900 name_to_id: dict[str, str] = {}
1901 missing_status: dict[str, list[str]] = {}
1902 for id_, job in jobs.items():
1903 if job["DAGNodeName"].startswith("wms_"):
1904 name_to_id[job["DAGNodeName"]] = id_
1905 if "NodeStatus" not in job or job["NodeStatus"] == NodeStatus.NOT_READY:
1906 missing_status.setdefault(job["from_dag_job"], []).append(id_)
1908 for name, dag_id in name_to_id.items():
1909 dag_status = jobs[dag_id].get("NodeStatus", NodeStatus.NOT_READY)
1910 if dag_status in {NodeStatus.NOT_READY, NodeStatus.FUTILE}:
1911 for id_ in missing_status.get(name, []):
1912 jobs[id_]["NodeStatus"] = dag_status
1914 return jobs
1917def read_single_dag_log(log_filename: str | os.PathLike) -> tuple[str, dict[str, dict[str, Any]]]:
1918 """Read job information from the DAGMan log file.
1920 Parameters
1921 ----------
1922 log_filename : `str` or `os.PathLike`
1923 DAGMan log filename.
1925 Returns
1926 -------
1927 wms_workflow_id : `str`
1928 HTCondor job id (i.e., <ClusterId>.<ProcId>) of the DAGMan job.
1929 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1930 HTCondor job information read from the log file mapped to HTCondor
1931 job id.
1933 Raises
1934 ------
1935 FileNotFoundError
1936 If cannot find DAGMan log in given wms_path.
1937 """
1938 wms_workflow_id = MISSING_ID
1939 dag_info: dict[str, dict[str, Any]] = {}
1941 filename = Path(log_filename)
1942 if filename.exists():
1943 _LOG.debug("dag node log filename: %s", filename)
1945 info: dict[str, Any] = {}
1946 job_event_log = htcondor.JobEventLog(str(filename))
1947 for event in job_event_log.events(stop_after=0):
1948 id_ = f"{event['Cluster']}.{event['Proc']}"
1949 if id_ not in info:
1950 info[id_] = {}
1951 wms_workflow_id = id_ # taking last job id in case of restarts
1952 _update_dicts(info[id_], event)
1953 info[id_][f"{event.type.name.lower()}_time"] = event["EventTime"]
1955 # only save latest DAG job
1956 dag_info = {wms_workflow_id: info[wms_workflow_id]}
1958 return wms_workflow_id, dag_info
1961def read_dag_log(wms_path: str | os.PathLike) -> tuple[str, dict[str, Any]]:
1962 """Read job information from the DAGMan log file.
1964 Parameters
1965 ----------
1966 wms_path : `str` or `os.PathLike`
1967 Path containing the DAGMan log file.
1969 Returns
1970 -------
1971 wms_workflow_id : `str`
1972 HTCondor job id (i.e., <ClusterId>.<ProcId>) of the DAGMan job.
1973 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1974 HTCondor job information read from the log file mapped to HTCondor
1975 job id.
1977 Raises
1978 ------
1979 FileNotFoundError
1980 If cannot find DAGMan log in given wms_path.
1981 """
1982 wms_workflow_id = MISSING_ID
1983 dag_info: dict[str, dict[str, Any]] = {}
1985 path = Path(wms_path)
1986 if path.exists(): 1986 ↛ 1998line 1986 didn't jump to line 1998 because the condition on line 1986 was always true
1987 # Can be more than one dag file in directory. Assume one
1988 # with lowest ID is main DAG
1989 ids = []
1990 for filename in path.glob("*.dag.dagman.log"):
1991 _LOG.debug("dag log filename: %s", filename)
1992 single_id, single_dag_info = read_single_dag_log(filename)
1993 _update_dicts(dag_info, single_dag_info)
1994 ids.append(single_id)
1995 if ids:
1996 wms_workflow_id = min(ids)
1998 if wms_workflow_id == MISSING_ID:
1999 raise FileNotFoundError(f"DAGMan log not found in {wms_path}")
2001 return wms_workflow_id, dag_info
2004def read_single_dag_nodes_log(filename: str | os.PathLike) -> dict[str, dict[str, Any]]:
2005 """Read job information from the DAGMan nodes log file.
2007 Parameters
2008 ----------
2009 filename : `str` or `os.PathLike`
2010 Path containing the DAGMan nodes log file.
2012 Returns
2013 -------
2014 info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2015 HTCondor job information read from the log file mapped to HTCondor
2016 job id.
2018 Raises
2019 ------
2020 FileNotFoundError
2021 If cannot find DAGMan node log in given wms_path.
2022 """
2023 _LOG.debug("dag node log filename: %s", filename)
2024 filename = Path(filename)
2026 info: dict[str, dict[str, Any]] = {}
2027 if not filename.exists():
2028 raise FileNotFoundError(f"{filename} does not exist")
2030 try:
2031 job_event_log = htcondor.JobEventLog(str(filename))
2032 except htcondor.HTCondorIOError as ex:
2033 _LOG.error("Problem reading nodes log file (%s): %s", filename, ex)
2034 import traceback
2036 traceback.print_stack()
2037 raise
2038 for event in job_event_log.events(stop_after=0):
2039 _LOG.debug("log event type = %s, keys = %s", event["EventTypeNumber"], event.keys())
2041 try:
2042 id_ = f"{event['Cluster']}.{event['Proc']}"
2043 except KeyError:
2044 _LOG.warn(
2045 "Log event missing ids (DAGNodeName=%s, EventTime=%s, EventTypeNumber=%s)",
2046 event.get("DAGNodeName", "UNK"),
2047 event.get("EventTime", "UNK"),
2048 event.get("EventTypeNumber", "UNK"),
2049 )
2050 else:
2051 if id_ not in info:
2052 info[id_] = {}
2053 # Workaround: Please check to see if still problem in
2054 # future HTCondor versions. Sometimes get a
2055 # JobAbortedEvent for a subdag job after it already
2056 # terminated normally. Seems to happen when using job
2057 # plus subdags.
2058 if event["EventTypeNumber"] == 9 and info[id_].get("EventTypeNumber", -1) == 5: 2058 ↛ 2059line 2058 didn't jump to line 2059 because the condition on line 2058 was never true
2059 _LOG.debug("Skipping spurious JobAbortedEvent: %s", dict(event))
2060 elif event["EventTypeNumber"] == 16 and event["DAGNodeName"] == "finalJob":
2061 # FINAL job's post script exit code is special and indicates
2062 # status of DAG instead of just the FINAL job. Save the
2063 # information separately.
2064 info[id_]["post"] = dict(event)
2065 else:
2066 _update_dicts(info[id_], event)
2067 info[id_][f"{event.type.name.lower()}_time"] = event["EventTime"]
2069 return info
2072def read_dag_nodes_log(wms_path: str | os.PathLike) -> dict[str, dict[str, Any]]:
2073 """Read job information from the DAGMan nodes log file.
2075 Parameters
2076 ----------
2077 wms_path : `str` or `os.PathLike`
2078 Path containing the DAGMan nodes log file.
2080 Returns
2081 -------
2082 info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2083 HTCondor job information read from the log file mapped to HTCondor
2084 job id.
2086 Raises
2087 ------
2088 FileNotFoundError
2089 If cannot find DAGMan node log in given wms_path.
2090 """
2091 info: dict[str, dict[str, Any]] = {}
2092 for filename in Path(wms_path).glob("*.dag.nodes.log"):
2093 _LOG.debug("dag node log filename: %s", filename)
2094 _update_dicts(info, read_single_dag_nodes_log(filename))
2096 # If submitted, the main nodes log file should exist
2097 if not info:
2098 raise FileNotFoundError(f"DAGMan node log not found in {wms_path}")
2100 # Subdags will not have dag nodes log files if they haven't
2101 # started running yet (so missing is not an error).
2102 for filename in Path(wms_path).glob("subdags/*/*.dag.nodes.log"):
2103 _LOG.debug("dag node log filename: %s", filename)
2104 _update_dicts(info, read_single_dag_nodes_log(filename))
2106 return info
2109def read_dag_info(wms_path: str | os.PathLike) -> tuple[Path, dict[str, dict[str, Any]]]:
2110 """Read custom DAGMan job information from the file.
2112 Parameters
2113 ----------
2114 wms_path : `str` or `os.PathLike`
2115 Path containing the file with the DAGMan job info.
2117 Returns
2118 -------
2119 filename : `pathlib.Path`
2120 Name of file containing the dag information.
2121 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2122 HTCondor job information.
2124 Raises
2125 ------
2126 FileNotFoundError
2127 If cannot find DAGMan job info file in the given location.
2128 """
2129 dag_info: dict[str, dict[str, Any]] = {}
2130 try:
2131 filename = next(Path(wms_path).glob("*.info.json"))
2132 except StopIteration as exc:
2133 raise FileNotFoundError(f"File with DAGMan job information not found in {wms_path}") from exc
2134 _LOG.debug("DAGMan job information filename: %s", filename)
2135 try:
2136 with open(filename) as fh:
2137 dag_info = json.load(fh)
2138 except (OSError, PermissionError) as exc:
2139 _LOG.debug("Retrieving DAGMan job information failed: %s", exc)
2140 return filename, dag_info
2143def write_dag_info(filename: str, schedd_dag_info: dict[str, dict[str, Any]]):
2144 """Write custom job information about DAGMan job.
2146 Parameters
2147 ----------
2148 filename : `str`
2149 Name of the file where the information will be stored. If not given,
2150 creates a filename using bps_run value.
2151 schedd_dag_info : `dict` [`str` `dict` [`str`, `~typing.Any`]]
2152 Information about the DAGMan job.
2153 """
2154 _LOG.debug("schedd_dag_info = %s", schedd_dag_info)
2155 schedd_name, dag_info = next(iter(schedd_dag_info.items()))
2156 dag_id, dag_ad = next(iter(dag_info.items()))
2157 ad = {"ClusterId": dag_ad["ClusterId"], "GlobalJobId": dag_ad["GlobalJobId"]}
2158 ad.update({key: val for key, val in dag_ad.items() if key.startswith("bps")})
2159 try:
2160 with open(filename, "w") as fh:
2161 info = {schedd_name: {dag_id: ad}}
2162 json.dump(info, fh)
2163 except (KeyError, OSError, PermissionError) as exc:
2164 _LOG.debug("Persisting DAGMan job information failed: %s", exc)
2166 return filename
2169def htc_tweak_log_info(wms_path: str | Path, job: dict[str, Any]) -> None:
2170 """Massage the given job info has same structure as if came from condor_q.
2172 Parameters
2173 ----------
2174 wms_path : `str` | `os.PathLike`
2175 Path containing an HTCondor event log file.
2176 job : `dict` [ `str`, `~typing.Any` ]
2177 A mapping between HTCondor job id and job information read from
2178 the log.
2179 """
2180 _LOG.debug("htc_tweak_log_info: %s %s", wms_path, job)
2182 # Use the presence of 'MyType' key as a proxy to determine if the job ad
2183 # contains the info extracted from the event log. Exit early if it doesn't
2184 # (e.g. it is a job ad for a pruned job).
2185 if "MyType" not in job:
2186 return
2188 try:
2189 job["ClusterId"] = job["Cluster"]
2190 job["ProcId"] = job["Proc"]
2191 except KeyError as e:
2192 _LOG.error("Missing key %s in job: %s", str(e), job)
2193 raise
2194 job["Iwd"] = str(wms_path)
2195 job["Owner"] = Path(wms_path).owner()
2197 if "LogNotes" in job:
2198 m = re.match(r"DAG Node: (\S+)", job["LogNotes"])
2199 if m: 2199 ↛ 2202line 2199 didn't jump to line 2202 because the condition on line 2199 was always true
2200 job["DAGNodeName"] = m.group(1)
2202 match job["MyType"]:
2203 case "ExecuteEvent":
2204 job["JobStatus"] = htcondor.JobStatus.RUNNING
2205 case "JobTerminatedEvent" | "PostScriptTerminatedEvent":
2206 job["JobStatus"] = htcondor.JobStatus.COMPLETED
2207 case "SubmitEvent":
2208 job["JobStatus"] = htcondor.JobStatus.IDLE
2209 case "JobAbortedEvent":
2210 job["JobStatus"] = htcondor.JobStatus.REMOVED
2211 case "JobHeldEvent":
2212 job["JobStatus"] = htcondor.JobStatus.HELD
2213 case "JobReleaseEvent":
2214 # If the job managing the execution of the root DAG is held and
2215 # released this will be the last event showing up in its
2216 # job event log even if the job is still running. If this is
2217 # the last event for a job corresponding to the workflow node
2218 # (either a normal payload job or the job managing the execution
2219 # of an inner DAG), its final status will be determined later
2220 # using node status log (see _htc_status_to_wms_state()).
2221 job["JobStatus"] = htcondor.JobStatus.RUNNING if "DAGNodeName" not in job else None
2222 case _:
2223 _LOG.debug("Unknown log event type: %s", job["MyType"])
2224 job["JobStatus"] = None
2226 # Use available information to add either "ExitCode" or "ExitSignal"
2227 # attribute that captures respectively job's exit status (if it finished
2228 # on its own accord) or its exit signal (if it was terminated by
2229 # a signal). Also, include a flag "ExitBySignal" to make distinguishing
2230 # between these two cases easy later on.
2231 if job["JobStatus"] in {
2232 htcondor.JobStatus.COMPLETED,
2233 htcondor.JobStatus.HELD,
2234 htcondor.JobStatus.REMOVED,
2235 }:
2236 new_job = HTC_JOB_AD_HANDLERS.handle(job)
2237 if new_job is not None:
2238 job = new_job
2239 else:
2240 _LOG.error("Could not determine exit status for job '%s.%s'", job["ClusterId"], job["ProcId"])
2243def htc_check_dagman_output(wms_path: str | os.PathLike) -> str:
2244 """Check the DAGMan output for error messages.
2246 Parameters
2247 ----------
2248 wms_path : `str` or `os.PathLike`
2249 Directory containing the DAGman output file.
2251 Returns
2252 -------
2253 message : `str`
2254 Message containing error messages from the DAGMan output. Empty
2255 string if no messages.
2257 Raises
2258 ------
2259 FileNotFoundError
2260 If cannot find DAGMan standard output file in given wms_path.
2261 """
2262 try:
2263 filename = next(Path(wms_path).glob("*.dag.dagman.out"))
2264 except StopIteration as exc:
2265 raise FileNotFoundError(f"DAGMan standard output file not found in {wms_path}") from exc
2266 _LOG.debug("dag output filename: %s", filename)
2268 p = re.compile(r"^(\d\d/\d\d/\d\d \d\d:\d\d:\d\d) (Job submit try \d+/\d+ failed|Warning:.*$|ERROR:.*$)")
2270 message = ""
2271 # Shouldn't normally see this warning. Initializing it just in case.
2272 last_warning = (
2273 "Missing warning: Usually means the .dagman.out file has been truncated "
2274 "or temporarily had a full quota."
2275 )
2276 try:
2277 with open(filename) as fh:
2278 last_submit_failed = "" # Since submit retries multiple times only report last one
2279 for line in fh:
2280 m = p.match(line)
2281 if m:
2282 if m.group(2).startswith("Job submit try"):
2283 last_submit_failed = m.group(1)
2284 elif m.group(2).startswith("ERROR: submit attempt failed"):
2285 pass # Should be handled by Job submit try
2286 elif m.group(2).startswith("Warning"):
2287 if ".dag.nodes.log is in /tmp" in m.group(2):
2288 last_warning = "Cannot submit from /tmp."
2289 else:
2290 last_warning = m.group(2)
2291 elif m.group(2) == "ERROR: Warning is fatal error because of DAGMAN_USE_STRICT setting":
2292 message += "ERROR: "
2293 message += last_warning
2294 message += "\n"
2295 elif m.group(2) in [ 2295 ↛ 2301line 2295 didn't jump to line 2301 because the condition on line 2295 was always true
2296 "ERROR: the following job(s) failed:",
2297 "ERROR: the following Node(s) failed:",
2298 ]:
2299 pass
2300 else:
2301 message += m.group(2)
2302 message += "\n"
2304 if last_submit_failed:
2305 message += f"Warn: Job submission issues (last: {last_submit_failed})"
2306 except (OSError, PermissionError):
2307 message = f"Warn: Could not read dagman output file from {wms_path}."
2308 _LOG.debug("dag output file message: %s", message)
2309 return message
2312def _read_rescue_headers(infh: TextIO) -> list[str]:
2313 """Read header lines from a rescue file.
2315 Parameters
2316 ----------
2317 infh : `TextIO`
2318 The rescue file from which to read the header lines.
2320 Returns
2321 -------
2322 header_lines : `list` [`str`]
2323 Header lines read from the rescue file.
2324 """
2325 header_lines = []
2326 for line in infh:
2327 line = line.strip()
2328 if not line.startswith("#"):
2329 break
2330 header_lines.append(line)
2331 return header_lines
2334def _update_rescue_headers(header_lines: list[str]) -> list[str]:
2335 """Update rescue header lines in place.
2337 Replaces ``wms_check_status`` node name prefixes with the corresponding
2338 group node names and adjusts the count of nodes premarked DONE.
2340 Parameters
2341 ----------
2342 header_lines : `list` [`str`]
2343 Header lines to update in place.
2345 Returns
2346 -------
2347 failed_subdags : `list` [`str`]
2348 Names of failed subdag jobs.
2349 """
2350 failed_subdags = []
2352 # Relies on HTCondor writing all failed node names on a single line
2353 # (see src/condor_dagman/dag.cpp in https://github.com/htcondor/htcondor).
2354 for i, line in enumerate(header_lines):
2355 if line.startswith("# Nodes that failed:"):
2356 nodes = header_lines[i + 1][1:].strip().split(",")
2357 new_nodes = []
2358 for node in nodes:
2359 if node.startswith("wms_check_status"):
2360 node = node[17:]
2361 failed_subdags.append(node)
2362 new_nodes.append(node)
2363 header_lines[i + 1] = f"# {','.join(new_nodes)}"
2364 break
2366 for i, line in enumerate(header_lines):
2367 m = re.match(r"^(# Nodes premarked DONE:)\s+(\d+)", line)
2368 if m:
2369 header_lines[i] = f"{m.group(1)} {int(m.group(2)) - len(failed_subdags)}"
2370 break
2372 return failed_subdags
2375def _write_rescue_headers(header_lines: list[str], outfh: TextIO) -> None:
2376 """Write the header lines to the new rescue file.
2378 Parameters
2379 ----------
2380 header_lines : `list` [`str`]
2381 Header lines to write to the new rescue file.
2382 outfh : `TextIO`
2383 New rescue file.
2384 """
2385 for line in header_lines:
2386 print(line, file=outfh)
2387 print("", file=outfh)
2390def _copy_done_lines(failed_subdags: list[str], infh: TextIO, outfh: TextIO) -> None:
2391 """Copy the DONE lines from the original rescue file skipping
2392 the failed group jobs.
2394 Parameters
2395 ----------
2396 failed_subdags : `list` [`str`]
2397 List of job names for the failed subdags
2398 infh : `TextIO`
2399 Original rescue file to copy from.
2400 outfh : `TextIO`
2401 New rescue file to copy to.
2402 """
2403 for line in infh:
2404 line = line.strip()
2405 try:
2406 _, node_name = line.split()
2407 except ValueError:
2408 _LOG.error(f"Unexpected line in rescue file = '{line}'")
2409 raise
2410 if node_name not in failed_subdags:
2411 print(line, file=outfh)
2414def _update_rescue_file(rescue_file: Path) -> list[str]:
2415 """Update the subdag failures in the main rescue file.
2417 Parameters
2418 ----------
2419 rescue_file : `pathlib.Path`
2420 The main rescue file that needs to be updated.
2422 Returns
2423 -------
2424 failed_subdags : `list` [`str`]
2425 Names of failed subdag jobs found in the DAGMan rescue file.
2426 """
2427 # To reduce memory requirements, not reading entire file into memory.
2428 rescue_tmp = rescue_file.with_suffix(rescue_file.suffix + ".tmp")
2429 with open(rescue_file) as infh:
2430 header_lines = _read_rescue_headers(infh)
2431 failed_subdags = _update_rescue_headers(header_lines)
2432 with open(rescue_tmp, "w") as outfh:
2433 _write_rescue_headers(header_lines, outfh)
2434 _copy_done_lines(failed_subdags, infh, outfh)
2435 rescue_file.unlink()
2436 rescue_tmp.rename(rescue_file)
2437 return failed_subdags
2440def _update_dicts(dict1, dict2):
2441 """Update dict1 with info in dict2.
2443 (Basically an update for nested dictionaries.)
2445 Parameters
2446 ----------
2447 dict1 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2448 HTCondor job information to be updated.
2449 dict2 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2450 Additional HTCondor job information.
2451 """
2452 for key, value in dict2.items():
2453 if key in dict1 and isinstance(dict1[key], dict) and isinstance(value, dict): 2453 ↛ 2454line 2453 didn't jump to line 2454 because the condition on line 2453 was never true
2454 _update_dicts(dict1[key], value)
2455 else:
2456 dict1[key] = value
2459def _locate_schedds(locate_all=False):
2460 """Find out Scheduler daemons in an HTCondor pool.
2462 Parameters
2463 ----------
2464 locate_all : `bool`, optional
2465 If True, all available schedulers in the HTCondor pool will be located.
2466 False by default which means that the search will be limited to looking
2467 for the Scheduler running on a local host.
2469 Returns
2470 -------
2471 schedds : `dict` [`str`, `htcondor.Schedd`]
2472 A mapping between Scheduler names and Python objects allowing for
2473 interacting with them.
2474 """
2475 coll = htcondor.Collector()
2477 schedd_ads = []
2478 if locate_all:
2479 schedd_ads.extend(coll.locateAll(htcondor.DaemonTypes.Schedd))
2480 else:
2481 schedd_ads.append(coll.locate(htcondor.DaemonTypes.Schedd))
2482 return {ad["Name"]: htcondor.Schedd(ad) for ad in schedd_ads}