Coverage for python/lsst/ctrl/bps/htcondor/lssthtc.py: 79%
951 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-14 09:36 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-14 09:36 +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 ↛ 1591line 1539 didn't jump to line 1591 because the condition on line 1539 was always true
1540 job_name = m.group("jobname")
1541 name_parts = job_name.split("_")
1543 label = ""
1544 if m.group("dir"):
1545 dir_match = re.search(r"jobs/([^\s/]+)", m.group("dir"))
1546 if dir_match:
1547 label = dir_match.group(1)
1548 else:
1549 _LOG.debug("Parse DAG: unparsed dir = %s", line)
1550 elif m.group("subfile"): 1550 ↛ 1557line 1550 didn't jump to line 1557 because the condition on line 1550 was always true
1551 subfile_match = re.search(r"jobs/([^\s/]+)", m.group("subfile"))
1552 if subfile_match:
1553 label = m.group("subfile").split("/")[1]
1554 else:
1555 label = pegasus_name_to_label(job_name)
1557 match m.group("command"):
1558 case "JOB":
1559 if m.group("noop"):
1560 job_type = WmsNodeType.NOOP
1561 # wms_noop_label
1562 label = name_parts[2]
1563 elif m.group("wms"):
1564 if name_parts[1] == "check": 1564 ↛ 1569line 1564 didn't jump to line 1569 because the condition on line 1564 was always true
1565 job_type = WmsNodeType.SUBDAG_CHECK
1566 # wms_check_status_wms_group_label
1567 label = name_parts[5]
1568 else:
1569 _LOG.warning(
1570 "Unexpected skipping of dag line due to unknown wms job: %s", line
1571 )
1572 else:
1573 job_type = WmsNodeType.PAYLOAD
1574 if label == "init": 1574 ↛ 1575line 1574 didn't jump to line 1575 because the condition on line 1574 was never true
1575 label = "pipetaskInit"
1576 counts[label] += 1
1577 case "FINAL":
1578 job_type = WmsNodeType.FINAL
1579 counts[label] += 1 # final counts a payload job.
1580 case "SERVICE":
1581 job_type = WmsNodeType.SERVICE
1582 case "SUBDAG EXTERNAL": 1582 ↛ 1586line 1582 didn't jump to line 1586 because the pattern on line 1582 always matched
1583 job_type = WmsNodeType.SUBDAG
1584 label = name_parts[2]
1586 job_name_to_label[job_name] = label
1587 job_name_to_type[job_name] = job_type
1588 else:
1589 # The line should, but didn't match the pattern above. Probably
1590 # problems with regex.
1591 _LOG.warning("Unexpected skipping of dag line: %s", line)
1593 return counts, job_name_to_label, job_name_to_type
1596def summarize_dag(dir_name: str) -> tuple[str, dict[str, str], dict[str, WmsNodeType]]:
1597 """Build bps_run_summary string from dag file.
1599 Parameters
1600 ----------
1601 dir_name : `str`
1602 Path that includes dag file for a run.
1604 Returns
1605 -------
1606 summary : `str`
1607 Semi-colon separated list of job labels and counts
1608 (Same format as saved in dag classad).
1609 job_name_to_label : `dict` [`str`, `str`]
1610 Mapping of job names to job labels.
1611 job_name_to_type : `dict` [`str`, `lsst.ctrl.bps.htcondor.WmsNodeType`]
1612 Mapping of job names to job types
1613 (e.g., payload, final, service).
1614 """
1615 # Later code depends upon insertion order
1616 counts: Counter[str] = Counter() # counts of payload jobs per label
1617 job_name_to_label: dict[str, str] = {}
1618 job_name_to_type: dict[str, WmsNodeType] = {}
1619 for filename in Path(dir_name).glob("*.dag"):
1620 single_counts, single_job_name_to_label, single_job_name_to_type = count_jobs_in_single_dag(filename)
1621 counts += single_counts
1622 _update_dicts(job_name_to_label, single_job_name_to_label)
1623 _update_dicts(job_name_to_type, single_job_name_to_type)
1625 for filename in Path(dir_name).glob("subdags/*/*.dag"):
1626 single_counts, single_job_name_to_label, single_job_name_to_type = count_jobs_in_single_dag(filename)
1627 counts += single_counts
1628 _update_dicts(job_name_to_label, single_job_name_to_label)
1629 _update_dicts(job_name_to_type, single_job_name_to_type)
1631 summary = ";".join([f"{name}:{counts[name]}" for name in counts])
1632 _LOG.debug("summarize_dag: %s %s %s", summary, job_name_to_label, job_name_to_type)
1633 return summary, job_name_to_label, job_name_to_type
1636def pegasus_name_to_label(name):
1637 """Convert pegasus job name to a label for the report.
1639 Parameters
1640 ----------
1641 name : `str`
1642 Name of job.
1644 Returns
1645 -------
1646 label : `str`
1647 Label for job.
1648 """
1649 label = "UNK"
1650 if name.startswith("create_dir") or name.startswith("stage_in") or name.startswith("stage_out"): 1650 ↛ 1651line 1650 didn't jump to line 1651 because the condition on line 1650 was never true
1651 label = "pegasus"
1652 else:
1653 m = re.match(r"pipetask_(\d+_)?([^_]+)", name)
1654 if m: 1654 ↛ 1655line 1654 didn't jump to line 1655 because the condition on line 1654 was never true
1655 label = m.group(2)
1656 if label == "init":
1657 label = "pipetaskInit"
1659 return label
1662def read_single_dag_status(filename: str | os.PathLike) -> dict[str, Any]:
1663 """Read the node status file for DAG summary information.
1665 Parameters
1666 ----------
1667 filename : `str` or `Path.pathlib`
1668 Node status filename.
1670 Returns
1671 -------
1672 dag_ad : `dict` [`str`, `~typing.Any`]
1673 DAG summary information.
1674 """
1675 dag_ad: dict[str, Any] = {}
1677 # While this is probably more up to date than dag classad, only read from
1678 # file if need to.
1679 try:
1680 node_stat_file = Path(filename)
1681 _LOG.debug("Reading Node Status File %s", node_stat_file)
1682 with open(node_stat_file) as infh:
1683 dag_ad = dict(classad.parseNext(infh)) # pylint: disable=E1101
1685 if not dag_ad: 1685 ↛ 1687line 1685 didn't jump to line 1687 because the condition on line 1685 was never true
1686 # Pegasus check here
1687 metrics_file = node_stat_file.with_suffix(".dag.metrics")
1688 if metrics_file.exists():
1689 with open(metrics_file) as infh:
1690 metrics = json.load(infh)
1691 dag_ad["NodesTotal"] = metrics.get("jobs", 0)
1692 dag_ad["NodesFailed"] = metrics.get("jobs_failed", 0)
1693 dag_ad["NodesDone"] = metrics.get("jobs_succeeded", 0)
1694 metrics_file = node_stat_file.with_suffix(".metrics")
1695 with open(metrics_file) as infh:
1696 metrics = json.load(infh)
1697 dag_ad["NodesTotal"] = metrics["wf_metrics"]["total_jobs"]
1698 except (OSError, PermissionError):
1699 pass
1701 _LOG.debug("read_dag_status: %s", dag_ad)
1702 return dag_ad
1705def read_dag_status(wms_path: str | os.PathLike) -> dict[str, Any]:
1706 """Read the node status file for DAG summary information.
1708 Parameters
1709 ----------
1710 wms_path : `str` or `os.PathLike`
1711 Path that includes node status file for a run.
1713 Returns
1714 -------
1715 dag_ad : `dict` [`str`, `~typing.Any`]
1716 DAG summary information, counts summed across any subdags.
1717 """
1718 dag_ads: dict[str, Any] = {}
1719 path = Path(wms_path)
1720 try:
1721 node_stat_file = next(path.glob("*.node_status"))
1722 except StopIteration as exc:
1723 raise FileNotFoundError(f"DAGMan node status not found in {wms_path}") from exc
1725 dag_ads = read_single_dag_status(node_stat_file)
1727 for node_stat_file in path.glob("subdags/*/*.node_status"):
1728 dag_ad = read_single_dag_status(node_stat_file)
1729 dag_ads["JobProcsHeld"] += dag_ad.get("JobProcsHeld", 0)
1730 dag_ads["NodesPost"] += dag_ad.get("NodesPost", 0)
1731 dag_ads["JobProcsIdle"] += dag_ad.get("JobProcsIdle", 0)
1732 dag_ads["NodesTotal"] += dag_ad.get("NodesTotal", 0)
1733 dag_ads["NodesFailed"] += dag_ad.get("NodesFailed", 0)
1734 dag_ads["NodesDone"] += dag_ad.get("NodesDone", 0)
1735 dag_ads["NodesQueued"] += dag_ad.get("NodesQueued", 0)
1736 dag_ads["NodesPre"] += dag_ad.get("NodesReady", 0)
1737 dag_ads["NodesFutile"] += dag_ad.get("NodesFutile", 0)
1738 dag_ads["NodesUnready"] += dag_ad.get("NodesUnready", 0)
1740 return dag_ads
1743def read_single_node_status(filename: str | os.PathLike, init_fake_id: int) -> dict[str, Any]:
1744 """Read entire node status file.
1746 Parameters
1747 ----------
1748 filename : `str` or `pathlib.Path`
1749 Node status filename.
1750 init_fake_id : `int`
1751 Initial fake id value.
1753 Returns
1754 -------
1755 jobs : `dict` [`str`, `~typing.Any`]
1756 DAG summary information compiled from the node status file combined
1757 with the information found in the node event log.
1759 Currently, if the same job attribute is found in both files, its value
1760 from the event log takes precedence over the value from the node status
1761 file.
1762 """
1763 filename = Path(filename)
1765 # Get jobid info from other places to fill in gaps in info from node_status
1766 _, job_name_to_label, job_name_to_type = count_jobs_in_single_dag(filename.with_suffix(".dag"))
1767 loginfo: dict[str, dict[str, Any]] = {}
1768 try:
1769 wms_workflow_id, _ = read_single_dag_log(filename.with_suffix(".dag.dagman.log"))
1770 loginfo = read_single_dag_nodes_log(filename.with_suffix(".dag.nodes.log"))
1771 except (OSError, PermissionError):
1772 pass
1774 job_name_to_id: dict[str, str] = {}
1775 _LOG.debug("loginfo = %s", loginfo)
1776 log_job_name_to_id: dict[str, str] = {}
1777 for job_id, job_info in loginfo.items():
1778 if "LogNotes" in job_info: 1778 ↛ 1777line 1778 didn't jump to line 1777 because the condition on line 1778 was always true
1779 m = re.match(r"DAG Node: (\S+)", job_info["LogNotes"])
1780 if m: 1780 ↛ 1777line 1780 didn't jump to line 1777 because the condition on line 1780 was always true
1781 job_name = m.group(1)
1782 log_job_name_to_id[job_name] = job_id
1783 job_info["DAGNodeName"] = job_name
1784 job_info["wms_node_type"] = job_name_to_type[job_name]
1785 job_info["bps_job_label"] = job_name_to_label[job_name]
1787 jobs = {}
1788 fake_id = init_fake_id # For nodes that do not yet have a job id, give fake one
1789 try:
1790 with open(filename) as fh:
1791 for ad in classad.parseAds(fh):
1792 match ad["Type"]:
1793 case "DagStatus":
1794 # Skip DAG summary.
1795 pass
1796 case "NodeStatus":
1797 job_name = ad["Node"]
1798 if job_name in job_name_to_label: 1798 ↛ 1800line 1798 didn't jump to line 1800 because the condition on line 1798 was always true
1799 job_label = job_name_to_label[job_name]
1800 elif "_" in job_name:
1801 job_label = job_name.split("_")[1]
1802 else:
1803 job_label = job_name
1805 job = dict(ad)
1806 if job_name in log_job_name_to_id:
1807 job_id = str(log_job_name_to_id[job_name])
1808 _update_dicts(job, loginfo[job_id])
1809 else:
1810 job_id = str(fake_id)
1811 job = dict(ad)
1812 fake_id -= 1
1813 jobs[job_id] = job
1814 job_name_to_id[job_name] = job_id
1816 # Make job info as if came from condor_q.
1817 job["ClusterId"] = int(float(job_id))
1818 job["DAGManJobID"] = wms_workflow_id
1819 job["DAGNodeName"] = job_name
1820 job["bps_job_label"] = job_label
1821 job["wms_node_type"] = job_name_to_type[job_name]
1823 case "StatusEnd": 1823 ↛ 1826line 1823 didn't jump to line 1826 because the pattern on line 1823 always matched
1824 # Skip node status file "epilog".
1825 pass
1826 case _:
1827 _LOG.debug(
1828 "Ignoring unknown classad type '%s' in the node status file '%s'",
1829 ad["Type"],
1830 filename,
1831 )
1832 except (OSError, PermissionError):
1833 pass
1835 # Check for missing jobs (e.g., submission failure or not submitted yet)
1836 # Use dag info to create job placeholders
1837 for name in set(job_name_to_label) - set(job_name_to_id):
1838 if name in log_job_name_to_id: # job was in nodes.log, but not node_status
1839 job_id = str(log_job_name_to_id[name])
1840 job = dict(loginfo[job_id])
1841 else:
1842 job_id = str(fake_id)
1843 fake_id -= 1
1844 job = {}
1845 job["NodeStatus"] = NodeStatus.NOT_READY
1847 job["ClusterId"] = int(float(job_id))
1848 job["ProcId"] = 0
1849 job["DAGManJobID"] = wms_workflow_id
1850 job["DAGNodeName"] = name
1851 job["bps_job_label"] = job_name_to_label[name]
1852 job["wms_node_type"] = job_name_to_type[name]
1853 jobs[f"{job['ClusterId']}.{job['ProcId']}"] = job
1855 for job_info in jobs.values():
1856 job_info["from_dag_job"] = f"wms_{filename.stem}"
1858 return jobs
1861def read_node_status(wms_path: str | os.PathLike) -> dict[str, dict[str, Any]]:
1862 """Read entire node status file.
1864 Parameters
1865 ----------
1866 wms_path : `str` or `os.PathLike`
1867 Path that includes node status file for a run.
1869 Returns
1870 -------
1871 jobs : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1872 DAG summary information compiled from the node status file combined
1873 with the information found in the node event log.
1875 Currently, if the same job attribute is found in both files, its value
1876 from the event log takes precedence over the value from the node status
1877 file.
1878 """
1879 jobs: dict[str, dict[str, Any]] = {}
1880 init_fake_id = -1
1882 # subdags may not have run so wouldn't have node_status file
1883 # use dag files and let read_single_node_status handle missing
1884 # node_status file.
1885 for dag_filename in Path(wms_path).glob("*.dag"):
1886 filename = dag_filename.with_suffix(".node_status")
1887 info = read_single_node_status(filename, init_fake_id)
1888 init_fake_id -= len(info)
1889 _update_dicts(jobs, info)
1891 for dag_filename in Path(wms_path).glob("subdags/*/*.dag"):
1892 filename = dag_filename.with_suffix(".node_status")
1893 info = read_single_node_status(filename, init_fake_id)
1894 init_fake_id -= len(info)
1895 _update_dicts(jobs, info)
1897 # Propagate pruned from subdags to jobs
1898 name_to_id: dict[str, str] = {}
1899 missing_status: dict[str, list[str]] = {}
1900 for id_, job in jobs.items():
1901 if job["DAGNodeName"].startswith("wms_"):
1902 name_to_id[job["DAGNodeName"]] = id_
1903 if "NodeStatus" not in job or job["NodeStatus"] == NodeStatus.NOT_READY:
1904 missing_status.setdefault(job["from_dag_job"], []).append(id_)
1906 for name, dag_id in name_to_id.items():
1907 dag_status = jobs[dag_id].get("NodeStatus", NodeStatus.NOT_READY)
1908 if dag_status in {NodeStatus.NOT_READY, NodeStatus.FUTILE}:
1909 for id_ in missing_status.get(name, []):
1910 jobs[id_]["NodeStatus"] = dag_status
1912 return jobs
1915def read_single_dag_log(log_filename: str | os.PathLike) -> tuple[str, dict[str, dict[str, Any]]]:
1916 """Read job information from the DAGMan log file.
1918 Parameters
1919 ----------
1920 log_filename : `str` or `os.PathLike`
1921 DAGMan log filename.
1923 Returns
1924 -------
1925 wms_workflow_id : `str`
1926 HTCondor job id (i.e., <ClusterId>.<ProcId>) of the DAGMan job.
1927 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1928 HTCondor job information read from the log file mapped to HTCondor
1929 job id.
1931 Raises
1932 ------
1933 FileNotFoundError
1934 If cannot find DAGMan log in given wms_path.
1935 """
1936 wms_workflow_id = "0"
1937 dag_info: dict[str, dict[str, Any]] = {}
1939 filename = Path(log_filename)
1940 if filename.exists():
1941 _LOG.debug("dag node log filename: %s", filename)
1943 info: dict[str, Any] = {}
1944 job_event_log = htcondor.JobEventLog(str(filename))
1945 for event in job_event_log.events(stop_after=0):
1946 id_ = f"{event['Cluster']}.{event['Proc']}"
1947 if id_ not in info:
1948 info[id_] = {}
1949 wms_workflow_id = id_ # taking last job id in case of restarts
1950 _update_dicts(info[id_], event)
1951 info[id_][f"{event.type.name.lower()}_time"] = event["EventTime"]
1953 # only save latest DAG job
1954 dag_info = {wms_workflow_id: info[wms_workflow_id]}
1956 return wms_workflow_id, dag_info
1959def read_dag_log(wms_path: str | os.PathLike) -> tuple[str, dict[str, Any]]:
1960 """Read job information from the DAGMan log file.
1962 Parameters
1963 ----------
1964 wms_path : `str` or `os.PathLike`
1965 Path containing the DAGMan log file.
1967 Returns
1968 -------
1969 wms_workflow_id : `str`
1970 HTCondor job id (i.e., <ClusterId>.<ProcId>) of the DAGMan job.
1971 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
1972 HTCondor job information read from the log file mapped to HTCondor
1973 job id.
1975 Raises
1976 ------
1977 FileNotFoundError
1978 If cannot find DAGMan log in given wms_path.
1979 """
1980 wms_workflow_id = MISSING_ID
1981 dag_info: dict[str, dict[str, Any]] = {}
1983 path = Path(wms_path)
1984 if path.exists(): 1984 ↛ 1996line 1984 didn't jump to line 1996 because the condition on line 1984 was always true
1985 # Can be more than one dag file in directory. Assume one
1986 # with lowest ID is main DAG
1987 ids = []
1988 for filename in path.glob("*.dag.dagman.log"):
1989 _LOG.debug("dag log filename: %s", filename)
1990 single_id, single_dag_info = read_single_dag_log(filename)
1991 _update_dicts(dag_info, single_dag_info)
1992 ids.append(single_id)
1993 if ids:
1994 wms_workflow_id = min(ids)
1996 if wms_workflow_id == MISSING_ID:
1997 raise FileNotFoundError(f"DAGMan log not found in {wms_path}")
1999 return wms_workflow_id, dag_info
2002def read_single_dag_nodes_log(filename: str | os.PathLike) -> dict[str, dict[str, Any]]:
2003 """Read job information from the DAGMan nodes log file.
2005 Parameters
2006 ----------
2007 filename : `str` or `os.PathLike`
2008 Path containing the DAGMan nodes log file.
2010 Returns
2011 -------
2012 info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2013 HTCondor job information read from the log file mapped to HTCondor
2014 job id.
2016 Raises
2017 ------
2018 FileNotFoundError
2019 If cannot find DAGMan node log in given wms_path.
2020 """
2021 _LOG.debug("dag node log filename: %s", filename)
2022 filename = Path(filename)
2024 info: dict[str, dict[str, Any]] = {}
2025 if not filename.exists():
2026 raise FileNotFoundError(f"{filename} does not exist")
2028 try:
2029 job_event_log = htcondor.JobEventLog(str(filename))
2030 except htcondor.HTCondorIOError as ex:
2031 _LOG.error("Problem reading nodes log file (%s): %s", filename, ex)
2032 import traceback
2034 traceback.print_stack()
2035 raise
2036 for event in job_event_log.events(stop_after=0):
2037 _LOG.debug("log event type = %s, keys = %s", event["EventTypeNumber"], event.keys())
2039 try:
2040 id_ = f"{event['Cluster']}.{event['Proc']}"
2041 except KeyError:
2042 _LOG.warn(
2043 "Log event missing ids (DAGNodeName=%s, EventTime=%s, EventTypeNumber=%s)",
2044 event.get("DAGNodeName", "UNK"),
2045 event.get("EventTime", "UNK"),
2046 event.get("EventTypeNumber", "UNK"),
2047 )
2048 else:
2049 if id_ not in info:
2050 info[id_] = {}
2051 # Workaround: Please check to see if still problem in
2052 # future HTCondor versions. Sometimes get a
2053 # JobAbortedEvent for a subdag job after it already
2054 # terminated normally. Seems to happen when using job
2055 # plus subdags.
2056 if event["EventTypeNumber"] == 9 and info[id_].get("EventTypeNumber", -1) == 5: 2056 ↛ 2057line 2056 didn't jump to line 2057 because the condition on line 2056 was never true
2057 _LOG.debug("Skipping spurious JobAbortedEvent: %s", dict(event))
2058 elif event["EventTypeNumber"] == 16 and event["DAGNodeName"] == "finalJob":
2059 # FINAL job's post script exit code is special and indicates
2060 # status of DAG instead of just the FINAL job. Save the
2061 # information separately.
2062 info[id_]["post"] = dict(event)
2063 else:
2064 _update_dicts(info[id_], event)
2065 info[id_][f"{event.type.name.lower()}_time"] = event["EventTime"]
2067 return info
2070def read_dag_nodes_log(wms_path: str | os.PathLike) -> dict[str, dict[str, Any]]:
2071 """Read job information from the DAGMan nodes log file.
2073 Parameters
2074 ----------
2075 wms_path : `str` or `os.PathLike`
2076 Path containing the DAGMan nodes log file.
2078 Returns
2079 -------
2080 info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2081 HTCondor job information read from the log file mapped to HTCondor
2082 job id.
2084 Raises
2085 ------
2086 FileNotFoundError
2087 If cannot find DAGMan node log in given wms_path.
2088 """
2089 info: dict[str, dict[str, Any]] = {}
2090 for filename in Path(wms_path).glob("*.dag.nodes.log"):
2091 _LOG.debug("dag node log filename: %s", filename)
2092 _update_dicts(info, read_single_dag_nodes_log(filename))
2094 # If submitted, the main nodes log file should exist
2095 if not info:
2096 raise FileNotFoundError(f"DAGMan node log not found in {wms_path}")
2098 # Subdags will not have dag nodes log files if they haven't
2099 # started running yet (so missing is not an error).
2100 for filename in Path(wms_path).glob("subdags/*/*.dag.nodes.log"):
2101 _LOG.debug("dag node log filename: %s", filename)
2102 _update_dicts(info, read_single_dag_nodes_log(filename))
2104 return info
2107def read_dag_info(wms_path: str | os.PathLike) -> tuple[Path, dict[str, dict[str, Any]]]:
2108 """Read custom DAGMan job information from the file.
2110 Parameters
2111 ----------
2112 wms_path : `str` or `os.PathLike`
2113 Path containing the file with the DAGMan job info.
2115 Returns
2116 -------
2117 filename : `pathlib.Path`
2118 Name of file containing the dag information.
2119 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2120 HTCondor job information.
2122 Raises
2123 ------
2124 FileNotFoundError
2125 If cannot find DAGMan job info file in the given location.
2126 """
2127 dag_info: dict[str, dict[str, Any]] = {}
2128 try:
2129 filename = next(Path(wms_path).glob("*.info.json"))
2130 except StopIteration as exc:
2131 raise FileNotFoundError(f"File with DAGMan job information not found in {wms_path}") from exc
2132 _LOG.debug("DAGMan job information filename: %s", filename)
2133 try:
2134 with open(filename) as fh:
2135 dag_info = json.load(fh)
2136 except (OSError, PermissionError) as exc:
2137 _LOG.debug("Retrieving DAGMan job information failed: %s", exc)
2138 return filename, dag_info
2141def write_dag_info(filename: str, schedd_dag_info: dict[str, dict[str, Any]]):
2142 """Write custom job information about DAGMan job.
2144 Parameters
2145 ----------
2146 filename : `str`
2147 Name of the file where the information will be stored. If not given,
2148 creates a filename using bps_run value.
2149 schedd_dag_info : `dict` [`str` `dict` [`str`, `~typing.Any`]]
2150 Information about the DAGMan job.
2151 """
2152 _LOG.debug("schedd_dag_info = %s", schedd_dag_info)
2153 schedd_name, dag_info = next(iter(schedd_dag_info.items()))
2154 dag_id, dag_ad = next(iter(dag_info.items()))
2155 ad = {"ClusterId": dag_ad["ClusterId"], "GlobalJobId": dag_ad["GlobalJobId"]}
2156 ad.update({key: val for key, val in dag_ad.items() if key.startswith("bps")})
2157 try:
2158 with open(filename, "w") as fh:
2159 info = {schedd_name: {dag_id: ad}}
2160 json.dump(info, fh)
2161 except (KeyError, OSError, PermissionError) as exc:
2162 _LOG.debug("Persisting DAGMan job information failed: %s", exc)
2164 return filename
2167def htc_tweak_log_info(wms_path: str | Path, job: dict[str, Any]) -> None:
2168 """Massage the given job info has same structure as if came from condor_q.
2170 Parameters
2171 ----------
2172 wms_path : `str` | `os.PathLike`
2173 Path containing an HTCondor event log file.
2174 job : `dict` [ `str`, `~typing.Any` ]
2175 A mapping between HTCondor job id and job information read from
2176 the log.
2177 """
2178 _LOG.debug("htc_tweak_log_info: %s %s", wms_path, job)
2180 # Use the presence of 'MyType' key as a proxy to determine if the job ad
2181 # contains the info extracted from the event log. Exit early if it doesn't
2182 # (e.g. it is a job ad for a pruned job).
2183 if "MyType" not in job:
2184 return
2186 try:
2187 job["ClusterId"] = job["Cluster"]
2188 job["ProcId"] = job["Proc"]
2189 except KeyError as e:
2190 _LOG.error("Missing key %s in job: %s", str(e), job)
2191 raise
2192 job["Iwd"] = str(wms_path)
2193 job["Owner"] = Path(wms_path).owner()
2195 if "LogNotes" in job:
2196 m = re.match(r"DAG Node: (\S+)", job["LogNotes"])
2197 if m: 2197 ↛ 2200line 2197 didn't jump to line 2200 because the condition on line 2197 was always true
2198 job["DAGNodeName"] = m.group(1)
2200 match job["MyType"]:
2201 case "ExecuteEvent":
2202 job["JobStatus"] = htcondor.JobStatus.RUNNING
2203 case "JobTerminatedEvent" | "PostScriptTerminatedEvent":
2204 job["JobStatus"] = htcondor.JobStatus.COMPLETED
2205 case "SubmitEvent":
2206 job["JobStatus"] = htcondor.JobStatus.IDLE
2207 case "JobAbortedEvent":
2208 job["JobStatus"] = htcondor.JobStatus.REMOVED
2209 case "JobHeldEvent":
2210 job["JobStatus"] = htcondor.JobStatus.HELD
2211 case "JobReleaseEvent":
2212 # If the job managing the execution of the root DAG is held and
2213 # released this will be the last event showing up in its
2214 # job event log even if the job is still running. If this is
2215 # the last event for a job corresponding to the workflow node
2216 # (either a normal payload job or the job managing the execution
2217 # of an inner DAG), its final status will be determined later
2218 # using node status log (see _htc_status_to_wms_state()).
2219 job["JobStatus"] = htcondor.JobStatus.RUNNING if "DAGNodeName" not in job else None
2220 case _:
2221 _LOG.debug("Unknown log event type: %s", job["MyType"])
2222 job["JobStatus"] = None
2224 # Use available information to add either "ExitCode" or "ExitSignal"
2225 # attribute that captures respectively job's exit status (if it finished
2226 # on its own accord) or its exit signal (if it was terminated by
2227 # a signal). Also, include a flag "ExitBySignal" to make distinguishing
2228 # between these two cases easy later on.
2229 if job["JobStatus"] in {
2230 htcondor.JobStatus.COMPLETED,
2231 htcondor.JobStatus.HELD,
2232 htcondor.JobStatus.REMOVED,
2233 }:
2234 new_job = HTC_JOB_AD_HANDLERS.handle(job)
2235 if new_job is not None:
2236 job = new_job
2237 else:
2238 _LOG.error("Could not determine exit status for job '%s.%s'", job["ClusterId"], job["ProcId"])
2241def htc_check_dagman_output(wms_path: str | os.PathLike) -> str:
2242 """Check the DAGMan output for error messages.
2244 Parameters
2245 ----------
2246 wms_path : `str` or `os.PathLike`
2247 Directory containing the DAGman output file.
2249 Returns
2250 -------
2251 message : `str`
2252 Message containing error messages from the DAGMan output. Empty
2253 string if no messages.
2255 Raises
2256 ------
2257 FileNotFoundError
2258 If cannot find DAGMan standard output file in given wms_path.
2259 """
2260 try:
2261 filename = next(Path(wms_path).glob("*.dag.dagman.out"))
2262 except StopIteration as exc:
2263 raise FileNotFoundError(f"DAGMan standard output file not found in {wms_path}") from exc
2264 _LOG.debug("dag output filename: %s", filename)
2266 p = re.compile(r"^(\d\d/\d\d/\d\d \d\d:\d\d:\d\d) (Job submit try \d+/\d+ failed|Warning:.*$|ERROR:.*$)")
2268 message = ""
2269 try:
2270 with open(filename) as fh:
2271 last_submit_failed = "" # Since submit retries multiple times only report last one
2272 for line in fh:
2273 m = p.match(line)
2274 if m:
2275 if m.group(2).startswith("Job submit try"):
2276 last_submit_failed = m.group(1)
2277 elif m.group(2).startswith("ERROR: submit attempt failed"):
2278 pass # Should be handled by Job submit try
2279 elif m.group(2).startswith("Warning"):
2280 if ".dag.nodes.log is in /tmp" in m.group(2):
2281 last_warning = "Cannot submit from /tmp."
2282 else:
2283 last_warning = m.group(2)
2284 elif m.group(2) == "ERROR: Warning is fatal error because of DAGMAN_USE_STRICT setting":
2285 message += "ERROR: "
2286 message += last_warning
2287 message += "\n"
2288 elif m.group(2) in [ 2288 ↛ 2294line 2288 didn't jump to line 2294 because the condition on line 2288 was always true
2289 "ERROR: the following job(s) failed:",
2290 "ERROR: the following Node(s) failed:",
2291 ]:
2292 pass
2293 else:
2294 message += m.group(2)
2295 message += "\n"
2297 if last_submit_failed:
2298 message += f"Warn: Job submission issues (last: {last_submit_failed})"
2299 except (OSError, PermissionError):
2300 message = f"Warn: Could not read dagman output file from {wms_path}."
2301 _LOG.debug("dag output file message: %s", message)
2302 return message
2305def _read_rescue_headers(infh: TextIO) -> list[str]:
2306 """Read header lines from a rescue file.
2308 Parameters
2309 ----------
2310 infh : `TextIO`
2311 The rescue file from which to read the header lines.
2313 Returns
2314 -------
2315 header_lines : `list` [`str`]
2316 Header lines read from the rescue file.
2317 """
2318 header_lines = []
2319 for line in infh:
2320 line = line.strip()
2321 if not line.startswith("#"):
2322 break
2323 header_lines.append(line)
2324 return header_lines
2327def _update_rescue_headers(header_lines: list[str]) -> list[str]:
2328 """Update rescue header lines in place.
2330 Replaces ``wms_check_status`` node name prefixes with the corresponding
2331 group node names and adjusts the count of nodes premarked DONE.
2333 Parameters
2334 ----------
2335 header_lines : `list` [`str`]
2336 Header lines to update in place.
2338 Returns
2339 -------
2340 failed_subdags : `list` [`str`]
2341 Names of failed subdag jobs.
2342 """
2343 failed_subdags = []
2345 # Relies on HTCondor writing all failed node names on a single line
2346 # (see src/condor_dagman/dag.cpp in https://github.com/htcondor/htcondor).
2347 for i, line in enumerate(header_lines):
2348 if line.startswith("# Nodes that failed:"):
2349 nodes = header_lines[i + 1][1:].strip().split(",")
2350 new_nodes = []
2351 for node in nodes:
2352 if node.startswith("wms_check_status"):
2353 node = node[17:]
2354 failed_subdags.append(node)
2355 new_nodes.append(node)
2356 header_lines[i + 1] = f"# {','.join(new_nodes)}"
2357 break
2359 for i, line in enumerate(header_lines):
2360 m = re.match(r"^(# Nodes premarked DONE:)\s+(\d+)", line)
2361 if m:
2362 header_lines[i] = f"{m.group(1)} {int(m.group(2)) - len(failed_subdags)}"
2363 break
2365 return failed_subdags
2368def _write_rescue_headers(header_lines: list[str], outfh: TextIO) -> None:
2369 """Write the header lines to the new rescue file.
2371 Parameters
2372 ----------
2373 header_lines : `list` [`str`]
2374 Header lines to write to the new rescue file.
2375 outfh : `TextIO`
2376 New rescue file.
2377 """
2378 for line in header_lines:
2379 print(line, file=outfh)
2380 print("", file=outfh)
2383def _copy_done_lines(failed_subdags: list[str], infh: TextIO, outfh: TextIO) -> None:
2384 """Copy the DONE lines from the original rescue file skipping
2385 the failed group jobs.
2387 Parameters
2388 ----------
2389 failed_subdags : `list` [`str`]
2390 List of job names for the failed subdags
2391 infh : `TextIO`
2392 Original rescue file to copy from.
2393 outfh : `TextIO`
2394 New rescue file to copy to.
2395 """
2396 for line in infh:
2397 line = line.strip()
2398 try:
2399 _, node_name = line.split()
2400 except ValueError:
2401 _LOG.error(f"Unexpected line in rescue file = '{line}'")
2402 raise
2403 if node_name not in failed_subdags:
2404 print(line, file=outfh)
2407def _update_rescue_file(rescue_file: Path) -> list[str]:
2408 """Update the subdag failures in the main rescue file.
2410 Parameters
2411 ----------
2412 rescue_file : `pathlib.Path`
2413 The main rescue file that needs to be updated.
2415 Returns
2416 -------
2417 failed_subdags : `list` [`str`]
2418 Names of failed subdag jobs found in the DAGMan rescue file.
2419 """
2420 # To reduce memory requirements, not reading entire file into memory.
2421 rescue_tmp = rescue_file.with_suffix(rescue_file.suffix + ".tmp")
2422 with open(rescue_file) as infh:
2423 header_lines = _read_rescue_headers(infh)
2424 failed_subdags = _update_rescue_headers(header_lines)
2425 with open(rescue_tmp, "w") as outfh:
2426 _write_rescue_headers(header_lines, outfh)
2427 _copy_done_lines(failed_subdags, infh, outfh)
2428 rescue_file.unlink()
2429 rescue_tmp.rename(rescue_file)
2430 return failed_subdags
2433def _update_dicts(dict1, dict2):
2434 """Update dict1 with info in dict2.
2436 (Basically an update for nested dictionaries.)
2438 Parameters
2439 ----------
2440 dict1 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2441 HTCondor job information to be updated.
2442 dict2 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
2443 Additional HTCondor job information.
2444 """
2445 for key, value in dict2.items():
2446 if key in dict1 and isinstance(dict1[key], dict) and isinstance(value, dict): 2446 ↛ 2447line 2446 didn't jump to line 2447 because the condition on line 2446 was never true
2447 _update_dicts(dict1[key], value)
2448 else:
2449 dict1[key] = value
2452def _locate_schedds(locate_all=False):
2453 """Find out Scheduler daemons in an HTCondor pool.
2455 Parameters
2456 ----------
2457 locate_all : `bool`, optional
2458 If True, all available schedulers in the HTCondor pool will be located.
2459 False by default which means that the search will be limited to looking
2460 for the Scheduler running on a local host.
2462 Returns
2463 -------
2464 schedds : `dict` [`str`, `htcondor.Schedd`]
2465 A mapping between Scheduler names and Python objects allowing for
2466 interacting with them.
2467 """
2468 coll = htcondor.Collector()
2470 schedd_ads = []
2471 if locate_all:
2472 schedd_ads.extend(coll.locateAll(htcondor.DaemonTypes.Schedd))
2473 else:
2474 schedd_ads.append(coll.locate(htcondor.DaemonTypes.Schedd))
2475 return {ad["Name"]: htcondor.Schedd(ad) for ad in schedd_ads}