Coverage for python/lsst/ctrl/bps/htcondor/report_utils.py: 61%
294 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-08-30 09:10 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-08-30 09: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"""Utility functions used for reporting."""
30import logging
31import os
32import re
33from pathlib import Path
34from typing import Any
36import htcondor
38from lsst.ctrl.bps import (
39 WmsJobReport,
40 WmsRunReport,
41 WmsSpecificInfo,
42 WmsStates,
43)
45from .common_utils import _htc_status_to_wms_state
46from .lssthtc import (
47 MISSING_ID,
48 WmsNodeType,
49 condor_search,
50 htc_check_dagman_output,
51 htc_tweak_log_info,
52 pegasus_name_to_label,
53 read_dag_info,
54 read_dag_log,
55 read_dag_status,
56 read_node_status,
57 summarize_dag,
58)
60_LOG = logging.getLogger(__name__)
63def _get_status_from_id(
64 wms_workflow_id: str, hist: float, schedds: dict[str, htcondor.Schedd]
65) -> tuple[WmsStates, str]:
66 """Gather run information using workflow id.
68 Parameters
69 ----------
70 wms_workflow_id : `str`
71 Limit to specific run based on id.
72 hist : `float`
73 Limit history search to this many days.
74 schedds : `dict` [ `str`, `htcondor.Schedd` ]
75 HTCondor schedulers which to query for job information. If empty
76 dictionary, all queries will be run against the local scheduler only.
78 Returns
79 -------
80 state : `lsst.ctrl.bps.WmsStates`
81 Status for the corresponding run.
82 message : `str`
83 Message with extra error information.
84 """
85 _LOG.debug("_get_status_from_id: id=%s, hist=%s, schedds=%s", wms_workflow_id, hist, schedds)
87 message = ""
89 # Collect information about the job by querying HTCondor schedd and
90 # HTCondor history.
91 schedd_dag_info = _get_info_from_schedd(wms_workflow_id, hist, schedds)
92 if len(schedd_dag_info) == 1:
93 schedd_name = next(iter(schedd_dag_info))
94 dag_id = next(iter(schedd_dag_info[schedd_name]))
95 dag_ad = schedd_dag_info[schedd_name][dag_id]
96 state = _htc_status_to_wms_state(dag_ad)
97 else:
98 state = WmsStates.UNKNOWN
99 message = f"DAGMan job {wms_workflow_id} not found in queue or history. Check id or try path."
100 return state, message
103def _get_status_from_path(wms_path: str | os.PathLike) -> tuple[WmsStates, str]:
104 """Gather run status from a given run directory.
106 Parameters
107 ----------
108 wms_path : `str` | `os.PathLike`
109 The directory containing the submit side files (e.g., HTCondor files).
111 Returns
112 -------
113 state : `lsst.ctrl.bps.WmsStates`
114 Status for the run.
115 message : `str`
116 Message to be printed.
117 """
118 wms_path = Path(wms_path).resolve()
119 message = ""
120 try:
121 wms_workflow_id, dag_ad = read_dag_log(wms_path)
122 except FileNotFoundError:
123 wms_workflow_id = MISSING_ID
124 message = f"DAGMan log not found in {wms_path}. Check path."
126 if wms_workflow_id == MISSING_ID:
127 state = WmsStates.UNKNOWN
128 else:
129 htc_tweak_log_info(wms_path, dag_ad[wms_workflow_id])
130 state = _htc_status_to_wms_state(dag_ad[wms_workflow_id])
132 return state, message
135def _report_from_path(wms_path):
136 """Gather run information from a given run directory.
138 Parameters
139 ----------
140 wms_path : `str`
141 The directory containing the submit side files (e.g., HTCondor files).
143 Returns
144 -------
145 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
146 Run information for the detailed report. The key is the HTCondor id
147 and the value is a collection of report information for that run.
148 message : `str`
149 Message to be printed with the summary report.
150 """
151 wms_workflow_id, jobs, message = _get_info_from_path(wms_path)
152 if wms_workflow_id == MISSING_ID:
153 run_reports = {}
154 else:
155 run_reports = _create_detailed_report_from_jobs(wms_workflow_id, jobs)
156 return run_reports, message
159def _report_from_id(wms_workflow_id, hist, schedds=None):
160 """Gather run information using workflow id.
162 Parameters
163 ----------
164 wms_workflow_id : `str`
165 Limit to specific run based on id.
166 hist : `float`
167 Limit history search to this many days.
168 schedds : `dict` [ `str`, `htcondor.Schedd` ], optional
169 HTCondor schedulers which to query for job information. If None
170 (default), all queries will be run against the local scheduler only.
172 Returns
173 -------
174 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
175 Run information for the detailed report. The key is the HTCondor id
176 and the value is a collection of report information for that run.
177 message : `str`
178 Message to be printed with the summary report.
179 """
180 messages = []
182 # Collect information about the job by querying HTCondor schedd and
183 # HTCondor history.
184 schedd_dag_info = _get_info_from_schedd(wms_workflow_id, hist, schedds)
185 if len(schedd_dag_info) == 1:
186 # Extract the DAG info without altering the results of the query.
187 schedd_name = next(iter(schedd_dag_info))
188 dag_id = next(iter(schedd_dag_info[schedd_name]))
189 dag_ad = schedd_dag_info[schedd_name][dag_id]
191 # If the provided workflow id does not correspond to the one extracted
192 # from the DAGMan log file in the submit directory, rerun the query
193 # with the id found in the file.
194 #
195 # This is to cover the situation in which the user provided the old job
196 # id of a restarted run.
197 try:
198 path_dag_id, _ = read_dag_log(dag_ad["Iwd"])
199 except FileNotFoundError as exc:
200 # At the moment missing DAGMan log is pretty much a fatal error.
201 # So empty the DAG info to finish early (see the if statement
202 # below).
203 schedd_dag_info.clear()
204 messages.append(f"Cannot create the report for '{dag_id}': {exc}")
205 else:
206 if path_dag_id != dag_id:
207 schedd_dag_info = _get_info_from_schedd(path_dag_id, hist, schedds)
208 messages.append(
209 f"WARNING: Found newer workflow executions in same submit directory as id '{dag_id}'. "
210 "This normally occurs when a run is restarted. The report shown is for the most "
211 f"recent status with run id '{path_dag_id}'"
212 )
214 if len(schedd_dag_info) == 0:
215 run_reports = {}
216 elif len(schedd_dag_info) == 1:
217 _, dag_info = schedd_dag_info.popitem()
218 dag_id, dag_ad = dag_info.popitem()
220 # Create a mapping between jobs and their classads. The keys will
221 # be of format 'ClusterId.ProcId'.
222 job_info = {dag_id: dag_ad}
224 # Find jobs (nodes) belonging to that DAGMan job.
225 job_constraint = f"DAGManJobId == {int(float(dag_id))}"
226 schedd_job_info = condor_search(constraint=job_constraint, hist=hist, schedds=schedds)
227 if schedd_job_info:
228 _, node_info = schedd_job_info.popitem()
229 job_info.update(node_info)
231 # Collect additional pieces of information about jobs using HTCondor
232 # files in the submission directory.
233 _, path_jobs, message = _get_info_from_path(dag_ad["Iwd"])
234 _update_jobs(job_info, path_jobs)
235 if message:
236 messages.append(message)
237 run_reports = _create_detailed_report_from_jobs(dag_id, job_info)
238 else:
239 ids = [ad["GlobalJobId"] for dag_info in schedd_dag_info.values() for ad in dag_info.values()]
240 message = (
241 f"More than one job matches id '{wms_workflow_id}', "
242 f"their global ids are: {', '.join(ids)}. Rerun with one of the global ids"
243 )
244 messages.append(message)
245 run_reports = {}
247 message = "\n".join(messages)
248 return run_reports, message
251def _get_info_from_schedd(
252 wms_workflow_id: str, hist: float, schedds: dict[str, htcondor.Schedd]
253) -> dict[str, dict[str, dict[str, Any]]]:
254 """Gather run information from HTCondor.
256 Parameters
257 ----------
258 wms_workflow_id : `str`
259 Limit to specific run based on id.
260 hist : `float`
261 Limit history search to this many days.
262 schedds : `dict` [ `str`, `htcondor.Schedd` ]
263 HTCondor schedulers which to query for job information. If empty
264 dictionary, all queries will be run against the local scheduler only.
266 Returns
267 -------
268 schedd_dag_info : `dict` [`str`, `dict` [`str`, `dict` [`str` Any]]]
269 Information about jobs satisfying the search criteria where for each
270 Scheduler, local HTCondor job ids are mapped to their respective
271 classads.
272 """
273 _LOG.debug("_get_info_from_schedd: id=%s, hist=%s, schedds=%s", wms_workflow_id, hist, schedds)
275 dag_constraint = 'regexp("dagman$", Cmd)'
276 try:
277 cluster_id = int(float(wms_workflow_id))
278 except ValueError:
279 dag_constraint += f' && GlobalJobId == "{wms_workflow_id}"'
280 else:
281 dag_constraint += f" && ClusterId == {cluster_id}"
283 # With the current implementation of the condor_* functions the query
284 # will always return only one match per Scheduler.
285 #
286 # Even in the highly unlikely situation where HTCondor history (which
287 # condor_search queries too) is long enough to have jobs from before
288 # the cluster ids were rolled over (and as a result there is more then
289 # one job with the same cluster id) they will not show up in
290 # the results.
291 schedd_dag_info = condor_search(constraint=dag_constraint, hist=hist, schedds=schedds)
292 return schedd_dag_info
295def _get_info_from_path(wms_path: str | os.PathLike) -> tuple[str, dict[str, dict[str, Any]], str]:
296 """Gather run information from a given run directory.
298 Parameters
299 ----------
300 wms_path : `str` or `os.PathLike`
301 Directory containing HTCondor files.
303 Returns
304 -------
305 wms_workflow_id : `str`
306 The run id which is a DAGman job id.
307 jobs : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
308 Information about jobs read from files in the given directory.
309 The key is the HTCondor id and the value is a dictionary of HTCondor
310 keys and values.
311 message : `str`
312 Message to be printed with the summary report.
313 """
314 # Ensure path is absolute, in particular for folks helping
315 # debug failures that need to dig around submit files.
316 wms_path = Path(wms_path).resolve()
318 messages = []
319 try:
320 wms_workflow_id, jobs = read_dag_log(wms_path)
321 _LOG.debug("_get_info_from_path: from dag log %s = %s", wms_workflow_id, jobs)
322 _update_jobs(jobs, read_node_status(wms_path))
323 _LOG.debug("_get_info_from_path: after node status %s = %s", wms_workflow_id, jobs)
325 # Add more info for DAGman job
326 job = jobs[wms_workflow_id]
327 job.update(read_dag_status(wms_path))
329 job["total_jobs"], job["state_counts"] = _get_state_counts_from_jobs(wms_workflow_id, jobs)
330 if "bps_run" not in job: 330 ↛ 333line 330 didn't jump to line 333 because the condition on line 330 was always true
331 _add_run_info(wms_path, job)
333 message = htc_check_dagman_output(wms_path)
334 if message: 334 ↛ 335line 334 didn't jump to line 335 because the condition on line 334 was never true
335 messages.append(message)
336 _LOG.debug(
337 "_get_info: id = %s, total_jobs = %s", wms_workflow_id, jobs[wms_workflow_id]["total_jobs"]
338 )
340 # Add extra pieces of information which cannot be found in HTCondor
341 # generated files like 'GlobalJobId'.
342 #
343 # Do not treat absence of this file as a serious error. Neither runs
344 # submitted with earlier versions of the plugin nor the runs submitted
345 # with Pegasus plugin will have it at the moment. However, once enough
346 # time passes and Pegasus plugin will have its own report() method
347 # (instead of sneakily using HTCondor's one), the lack of that file
348 # should be treated as seriously as lack of any other file.
349 try:
350 _, job_info = read_dag_info(wms_path)
351 except FileNotFoundError as exc:
352 message = f"Warn: Some information may not be available: {exc}"
353 messages.append(message)
354 else:
355 schedd_name = next(iter(job_info))
356 job_ad = next(iter(job_info[schedd_name].values()))
357 job.update(job_ad)
358 except FileNotFoundError as err:
359 message = f"Could not find HTCondor files in '{wms_path}' ({err})"
360 _LOG.debug(message)
361 messages.append(message)
362 try:
363 message = htc_check_dagman_output(wms_path)
364 except FileNotFoundError as err:
365 message = (
366 f"Could not find DAGMan standard output file in '{wms_path}'.\n"
367 "Check that path is a valid HTCondor submission directory.\n"
368 "Check that the condor_dagman executable path is accessible from the AP machine."
369 )
370 if message:
371 messages.append(message)
372 wms_workflow_id = MISSING_ID
373 jobs = {}
375 # Add more condor_q-like info.
376 for job in jobs.values():
377 htc_tweak_log_info(wms_path, job)
379 message = "\n".join([msg for msg in messages if msg])
380 _LOG.debug("wms_workflow_id = %s, jobs = %s", wms_workflow_id, jobs.keys())
381 _LOG.debug("message = %s", message)
382 return wms_workflow_id, jobs, message
385def _create_detailed_report_from_jobs(
386 wms_workflow_id: str, jobs: dict[str, dict[str, Any]]
387) -> dict[str, WmsRunReport]:
388 """Gather run information to be used in generating summary reports.
390 Parameters
391 ----------
392 wms_workflow_id : `str`
393 The run id to create the report for.
394 jobs : `dict` [`str`, `dict` [`str`, Any]]
395 Mapping HTCondor job id to job information.
397 Returns
398 -------
399 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
400 Run information for the detailed report. The key is the given HTCondor
401 id and the value is a collection of report information for that run.
402 """
403 _LOG.debug("_create_detailed_report: id = %s, job = %s", wms_workflow_id, jobs[wms_workflow_id])
405 dag_ad = jobs[wms_workflow_id]
407 report = WmsRunReport(
408 wms_id=f"{dag_ad['ClusterId']}.{dag_ad['ProcId']}",
409 global_wms_id=dag_ad.get("GlobalJobId", "MISS"),
410 path=dag_ad["Iwd"],
411 label=dag_ad.get("bps_job_label", "MISS"),
412 run=dag_ad.get("bps_run", "MISS"),
413 site=dag_ad.get("bps_runsite", ""),
414 project=dag_ad.get("bps_project", "MISS"),
415 campaign=dag_ad.get("bps_campaign", "MISS"),
416 payload=dag_ad.get("bps_payload", "MISS"),
417 operator=_get_owner(dag_ad),
418 run_summary=_get_run_summary(dag_ad),
419 state=_htc_status_to_wms_state(dag_ad),
420 total_number_jobs=0,
421 jobs=[],
422 job_state_counts=dict.fromkeys(WmsStates, 0),
423 exit_code_summary={},
424 )
426 payload_jobs = {} # keep track for later processing
427 specific_info = WmsSpecificInfo()
428 for job_id, job_ad in jobs.items():
429 if job_ad.get("wms_node_type", WmsNodeType.UNKNOWN) in [WmsNodeType.PAYLOAD, WmsNodeType.FINAL]:
430 try:
431 name = job_ad.get("DAGNodeName", job_id)
432 wms_state = _htc_status_to_wms_state(job_ad)
433 job_report = WmsJobReport(
434 wms_id=job_id,
435 name=name,
436 label=job_ad.get("bps_job_label", pegasus_name_to_label(name)),
437 state=wms_state,
438 )
439 if job_report.label == "init": 439 ↛ 440line 439 didn't jump to line 440 because the condition on line 439 was never true
440 job_report.label = "pipetaskInit"
441 report.job_state_counts[wms_state] += 1
442 report.jobs.append(job_report)
443 payload_jobs[job_id] = job_ad
444 except KeyError as ex:
445 _LOG.error("Job missing key '%s': %s", str(ex), job_ad)
446 raise
447 elif is_service_job(job_ad):
448 _LOG.debug(
449 "Found service job: id='%s', name='%s', label='%s', NodeStatus='%s', JobStatus='%s'",
450 job_id,
451 job_ad["DAGNodeName"],
452 job_ad.get("bps_job_label", "MISS"),
453 job_ad.get("NodeStatus", "MISS"),
454 job_ad.get("JobStatus", "MISS"),
455 )
456 _add_service_job_specific_info(job_ad, specific_info)
458 report.total_number_jobs = len(payload_jobs)
459 report.exit_code_summary = _get_exit_code_summary(payload_jobs)
460 if specific_info:
461 report.specific_info = specific_info
463 # Workflow will exit with non-zero DAG_STATUS if problem with
464 # any of the wms jobs. So change FAILED to SUCCEEDED if all
465 # payload jobs SUCCEEDED.
466 if report.total_number_jobs == report.job_state_counts[WmsStates.SUCCEEDED]:
467 report.state = WmsStates.SUCCEEDED
469 run_reports = {report.wms_id: report}
470 _LOG.debug("_create_detailed_report: run_reports = %s", run_reports)
471 return run_reports
474def _add_service_job_specific_info(job_ad: dict[str, Any], specific_info: WmsSpecificInfo) -> None:
475 """Generate report information for service job.
477 Parameters
478 ----------
479 job_ad : `dict` [`str`, `~typing.Any`]
480 Provisioning job information.
481 specific_info : `lsst.ctrl.bps.WmsSpecificInfo`
482 Where to add message.
483 """
484 status_details = ""
485 job_status = _htc_status_to_wms_state(job_ad)
487 # Service jobs in queue are deleted when DAG is done.
488 # To get accurate status, need to check other info.
489 if (
490 job_status == WmsStates.DELETED
491 and "Reason" in job_ad
492 and (
493 "Removed by DAGMan" in job_ad["Reason"]
494 or "removed because <OtherJobRemoveRequirements = DAGManJobId =?=" in job_ad["Reason"]
495 or "DAG is exiting and writing rescue file." in job_ad["Reason"]
496 )
497 ):
498 if "HoldReason" in job_ad:
499 # HoldReason exists even if released, so check.
500 if "job_released_time" in job_ad and job_ad["job_held_time"] < job_ad["job_released_time"]:
501 # If released, assume running until deleted.
502 job_status = WmsStates.SUCCEEDED
503 status_details = ""
504 else:
505 # If job held when deleted by DAGMan, still want to
506 # report hold reason
507 status_details = f"(Job was held for the following reason: {job_ad['HoldReason']})"
509 else:
510 job_status = WmsStates.SUCCEEDED
511 elif job_status == WmsStates.SUCCEEDED:
512 status_details = "(Note: Finished before workflow.)"
513 elif job_status == WmsStates.HELD:
514 status_details = f"({job_ad['HoldReason']})"
516 template = "Status of {job_name}: {status} {status_details}"
517 context = {
518 "job_name": job_ad["DAGNodeName"],
519 "status": job_status.name,
520 "status_details": status_details,
521 }
522 specific_info.add_message(template=template, context=context)
525def _summary_report(user, hist, pass_thru, schedds=None):
526 """Gather run information to be used in generating summary reports.
528 Parameters
529 ----------
530 user : `str`
531 Run lookup restricted to given user.
532 hist : `float`
533 How many previous days to search for run information.
534 pass_thru : `str`
535 Advanced users can define the HTCondor constraint to be used
536 when searching queue and history.
538 Returns
539 -------
540 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
541 Run information for the summary report. The keys are HTCondor ids and
542 the values are collections of report information for each run.
543 message : `str`
544 Message to be printed with the summary report.
545 """
546 # only doing summary report so only look for dagman jobs
547 if pass_thru:
548 constraint = pass_thru
549 else:
550 # Notes:
551 # * bps_isjob == 'True' isn't getting set for DAG jobs that are
552 # manually restarted.
553 # * Any job with DAGManJobID isn't a DAG job
554 constraint = 'bps_isjob == "True" && JobUniverse == 7 && DAGManJobID =?= Undefined'
555 if user:
556 constraint += f' && (Owner == "{user}" || bps_operator == "{user}")'
558 job_info = condor_search(constraint=constraint, hist=hist, schedds=schedds)
560 # Have list of DAGMan jobs, need to get run_report info.
561 run_reports = {}
562 msg = ""
563 for jobs in job_info.values():
564 for job_id, job in jobs.items():
565 total_jobs, state_counts = _get_state_counts_from_dag_job(job)
566 # If didn't get from queue information (e.g., Kerberos bug),
567 # try reading from file.
568 if total_jobs == 0:
569 try:
570 job.update(read_dag_status(job["Iwd"]))
571 total_jobs, state_counts = _get_state_counts_from_dag_job(job)
572 except (StopIteration, FileNotFoundError):
573 pass # don't kill report can't find htcondor files
575 if "bps_run" not in job:
576 _add_run_info(job["Iwd"], job)
577 report = WmsRunReport(
578 wms_id=job_id,
579 global_wms_id=job["GlobalJobId"],
580 path=job["Iwd"],
581 label=job.get("bps_job_label", "MISS"),
582 run=job.get("bps_run", "MISS"),
583 project=job.get("bps_project", "MISS"),
584 campaign=job.get("bps_campaign", "MISS"),
585 payload=job.get("bps_payload", "MISS"),
586 site=job.get("bps_runsite", ""),
587 operator=_get_owner(job),
588 run_summary=_get_run_summary(job),
589 state=_htc_status_to_wms_state(job),
590 jobs=[],
591 total_number_jobs=total_jobs,
592 job_state_counts=state_counts,
593 )
594 run_reports[report.global_wms_id] = report
596 return run_reports, msg
599def _add_run_info(wms_path, job):
600 """Find BPS run information elsewhere for runs without bps attributes.
602 Parameters
603 ----------
604 wms_path : `str`
605 Path to submit files for the run.
606 job : `dict` [`str`, `~typing.Any`]
607 HTCondor dag job information.
609 Raises
610 ------
611 StopIteration
612 If cannot find file it is looking for. Permission errors are
613 caught and job's run is marked with error.
614 """
615 path = Path(wms_path) / "jobs"
616 try:
617 subfile = next(path.glob("**/*.sub"))
618 except (StopIteration, PermissionError):
619 job["bps_run"] = "Unavailable"
620 else:
621 _LOG.debug("_add_run_info: subfile = %s", subfile)
622 try:
623 with open(subfile, encoding="utf-8") as fh:
624 for line in fh:
625 if line.startswith("+bps_"):
626 m = re.match(r"\+(bps_[^\s]+)\s*=\s*(.+)$", line)
627 if m:
628 _LOG.debug("Matching line: %s", line)
629 job[m.group(1)] = m.group(2).replace('"', "")
630 else:
631 _LOG.debug("Could not parse attribute: %s", line)
632 except PermissionError:
633 job["bps_run"] = "PermissionError"
634 _LOG.debug("After adding job = %s", job)
637def _get_owner(job):
638 """Get the owner of a dag job.
640 Parameters
641 ----------
642 job : `dict` [`str`, `~typing.Any`]
643 HTCondor dag job information.
645 Returns
646 -------
647 owner : `str`
648 Owner of the dag job.
649 """
650 owner = job.get("bps_operator", None)
651 if not owner: 651 ↛ 652line 651 didn't jump to line 652 because the condition on line 651 was never true
652 owner = job.get("Owner", None)
653 if not owner:
654 _LOG.warning("Could not get Owner from htcondor job: %s", job)
655 owner = "MISS"
656 return owner
659def _get_run_summary(job):
660 """Get the run summary for a job.
662 Parameters
663 ----------
664 job : `dict` [`str`, `~typing.Any`]
665 HTCondor dag job information.
667 Returns
668 -------
669 summary : `str`
670 Number of jobs per PipelineTask label in approximate pipeline order.
671 Format: <label>:<count>[;<label>:<count>]+
672 """
673 summary = job.get("bps_job_summary", job.get("bps_run_summary", None))
674 if not summary:
675 summary, _, _ = summarize_dag(job["Iwd"])
676 if not summary:
677 _LOG.warning("Could not get run summary for htcondor job: %s", job)
678 _LOG.debug("_get_run_summary: summary=%s", summary)
680 # Workaround sometimes using init vs pipetaskInit
681 summary = summary.replace("init:", "pipetaskInit:")
683 if "pegasus_version" in job and "pegasus" not in summary: 683 ↛ 684line 683 didn't jump to line 684 because the condition on line 683 was never true
684 summary += ";pegasus:0"
686 return summary
689def _get_exit_code_summary(jobs):
690 """Get the exit code summary for a run.
692 Parameters
693 ----------
694 jobs : `dict` [`str`, `dict` [`str`, Any]]
695 Mapping HTCondor job id to job information.
697 Returns
698 -------
699 summary : `dict` [`str`, `list` [`int`]]
700 Jobs' exit codes per job label.
701 """
702 summary = {}
703 for job_id, job_ad in jobs.items():
704 job_label = job_ad["bps_job_label"]
705 summary.setdefault(job_label, [])
706 try:
707 exit_code = 0
708 job_status = job_ad["JobStatus"]
709 match job_status:
710 case htcondor.JobStatus.COMPLETED | htcondor.JobStatus.HELD | htcondor.JobStatus.REMOVED:
711 exit_code = job_ad["ExitSignal"] if job_ad["ExitBySignal"] else job_ad["ExitCode"]
712 case (
713 htcondor.JobStatus.IDLE
714 | htcondor.JobStatus.RUNNING
715 | htcondor.JobStatus.TRANSFERRING_OUTPUT
716 | htcondor.JobStatus.SUSPENDED
717 ):
718 pass
719 case _:
720 _LOG.debug("Unknown 'JobStatus' value ('%d') in classad for job '%s'", job_status, job_id)
721 if exit_code != 0:
722 summary[job_label].append(exit_code)
723 except KeyError as ex:
724 _LOG.debug("Attribute '%s' not found in the classad for job '%s'", ex, job_id)
725 return summary
728def _get_state_counts_from_jobs(
729 wms_workflow_id: str, jobs: dict[str, dict[str, Any]]
730) -> tuple[int, dict[WmsStates, int]]:
731 """Count number of jobs per WMS state.
733 The workflow job and the service jobs are excluded from the count.
735 Parameters
736 ----------
737 wms_workflow_id : `str`
738 HTCondor job id.
739 jobs : `dict [`dict` [`str`, `~typing.Any`]]
740 HTCondor dag job information.
742 Returns
743 -------
744 total_count : `int`
745 Total number of dag nodes.
746 state_counts : `dict` [`lsst.ctrl.bps.WmsStates`, `int`]
747 Keys are the different WMS states and values are counts of jobs
748 that are in that WMS state.
749 """
750 state_counts = dict.fromkeys(WmsStates, 0)
751 for job_id, job_ad in jobs.items():
752 if job_id != wms_workflow_id and job_ad.get("wms_node_type", WmsNodeType.UNKNOWN) in [
753 WmsNodeType.PAYLOAD,
754 WmsNodeType.FINAL,
755 ]:
756 state_counts[_htc_status_to_wms_state(job_ad)] += 1
757 total_count = sum(state_counts.values())
759 return total_count, state_counts
762def _get_state_counts_from_dag_job(job):
763 """Count number of jobs per WMS state.
765 Parameters
766 ----------
767 job : `dict` [`str`, `~typing.Any`]
768 HTCondor dag job information.
770 Returns
771 -------
772 total_count : `int`
773 Total number of dag nodes.
774 state_counts : `dict` [`lsst.ctrl.bps.WmsStates`, `int`]
775 Keys are the different WMS states and values are counts of jobs
776 that are in that WMS state.
777 """
778 _LOG.debug("_get_state_counts_from_dag_job: job = %s %s", type(job), len(job))
779 state_counts = dict.fromkeys(WmsStates, 0)
780 if "DAG_NodesReady" in job: 780 ↛ 792line 780 didn't jump to line 792 because the condition on line 780 was always true
781 state_counts = {
782 WmsStates.UNREADY: job.get("DAG_NodesUnready", 0),
783 WmsStates.READY: job.get("DAG_NodesReady", 0),
784 WmsStates.HELD: job.get("DAG_JobsHeld", 0),
785 WmsStates.SUCCEEDED: job.get("DAG_NodesDone", 0),
786 WmsStates.FAILED: job.get("DAG_NodesFailed", 0),
787 WmsStates.PRUNED: job.get("DAG_NodesFutile", 0),
788 WmsStates.MISFIT: job.get("DAG_NodesPre", 0) + job.get("DAG_NodesPost", 0),
789 }
790 total_jobs = job.get("DAG_NodesTotal")
791 _LOG.debug("_get_state_counts_from_dag_job: from DAG_* keys, total_jobs = %s", total_jobs)
792 elif "NodesFailed" in job:
793 state_counts = {
794 WmsStates.UNREADY: job.get("NodesUnready", 0),
795 WmsStates.READY: job.get("NodesReady", 0),
796 WmsStates.HELD: job.get("JobProcsHeld", 0),
797 WmsStates.SUCCEEDED: job.get("NodesDone", 0),
798 WmsStates.FAILED: job.get("NodesFailed", 0),
799 WmsStates.PRUNED: job.get("NodesFutile", 0),
800 WmsStates.MISFIT: job.get("NodesPre", 0) + job.get("NodesPost", 0),
801 }
802 try:
803 total_jobs = job.get("NodesTotal")
804 except KeyError as ex:
805 _LOG.error("Job missing %s. job = %s", str(ex), job)
806 raise
807 _LOG.debug("_get_state_counts_from_dag_job: from NODES* keys, total_jobs = %s", total_jobs)
808 else:
809 # With Kerberos job auth and Kerberos bug, if warning would be printed
810 # for every DAG.
811 _LOG.debug("Can't get job state counts %s", job["Iwd"])
812 total_jobs = 0
814 _LOG.debug("total_jobs = %s, state_counts: %s", total_jobs, state_counts)
815 return total_jobs, state_counts
818def _update_jobs(jobs1, jobs2):
819 """Update jobs1 with info in jobs2.
821 (Basically an update for nested dictionaries.)
823 Parameters
824 ----------
825 jobs1 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
826 HTCondor job information to be updated.
827 jobs2 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
828 Additional HTCondor job information.
829 """
830 for job_id, job_ad in jobs2.items():
831 if job_id in jobs1: 831 ↛ 832line 831 didn't jump to line 832 because the condition on line 831 was never true
832 jobs1[job_id].update(job_ad)
833 else:
834 jobs1[job_id] = job_ad
837def is_service_job(job_ad: dict[str, Any]) -> bool:
838 """Determine if a job is a service one.
840 Parameters
841 ----------
842 job_ad : `dict` [`str`, Any]
843 Information about an HTCondor job.
845 Returns
846 -------
847 is_service_job : `bool`
848 True if the job is a service one, false otherwise.
850 Notes
851 -----
852 At the moment, HTCondor does not provide a native way to distinguish
853 between payload and service jobs in the workflow. This code depends
854 on read_node_status adding wms_node_type.
855 """
856 return job_ad.get("wms_node_type", WmsNodeType.UNKNOWN) == WmsNodeType.SERVICE