Coverage for python/lsst/ctrl/bps/htcondor/report_utils.py: 61%
295 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-29 02:24 -0700
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-29 02:24 -0700
1# This file is part of ctrl_bps_htcondor.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This software is dual licensed under the GNU General Public License and also
10# under a 3-clause BSD license. Recipients may choose which of these licenses
11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt,
12# respectively. If you choose the GPL option then the following text applies
13# (but note that there is still no warranty even if you opt for BSD instead):
14#
15# This program is free software: you can redistribute it and/or modify
16# it under the terms of the GNU General Public License as published by
17# the Free Software Foundation, either version 3 of the License, or
18# (at your option) any later version.
19#
20# This program is distributed in the hope that it will be useful,
21# but WITHOUT ANY WARRANTY; without even the implied warranty of
22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
23# GNU General Public License for more details.
24#
25# You should have received a copy of the GNU General Public License
26# along with this program. If not, see <https://www.gnu.org/licenses/>.
28"""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 dag_ad: dict[str, Any] = {}
121 try:
122 wms_workflow_id, dag_ad = read_dag_log(wms_path)
123 except FileNotFoundError:
124 wms_workflow_id = MISSING_ID
125 message = f"DAGMan log not found in {wms_path}. Check path."
127 if wms_workflow_id == MISSING_ID:
128 state = WmsStates.UNKNOWN
129 else:
130 htc_tweak_log_info(wms_path, dag_ad[wms_workflow_id])
131 state = _htc_status_to_wms_state(dag_ad[wms_workflow_id])
133 return state, message
136def _report_from_path(wms_path):
137 """Gather run information from a given run directory.
139 Parameters
140 ----------
141 wms_path : `str`
142 The directory containing the submit side files (e.g., HTCondor files).
144 Returns
145 -------
146 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
147 Run information for the detailed report. The key is the HTCondor id
148 and the value is a collection of report information for that run.
149 message : `str`
150 Message to be printed with the summary report.
151 """
152 wms_workflow_id, jobs, message = _get_info_from_path(wms_path)
153 if wms_workflow_id == MISSING_ID:
154 run_reports = {}
155 else:
156 run_reports = _create_detailed_report_from_jobs(wms_workflow_id, jobs)
157 return run_reports, message
160def _report_from_id(wms_workflow_id, hist, schedds=None):
161 """Gather run information using workflow id.
163 Parameters
164 ----------
165 wms_workflow_id : `str`
166 Limit to specific run based on id.
167 hist : `float`
168 Limit history search to this many days.
169 schedds : `dict` [ `str`, `htcondor.Schedd` ], optional
170 HTCondor schedulers which to query for job information. If None
171 (default), all queries will be run against the local scheduler only.
173 Returns
174 -------
175 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
176 Run information for the detailed report. The key is the HTCondor id
177 and the value is a collection of report information for that run.
178 message : `str`
179 Message to be printed with the summary report.
180 """
181 messages = []
183 # Collect information about the job by querying HTCondor schedd and
184 # HTCondor history.
185 schedd_dag_info = _get_info_from_schedd(wms_workflow_id, hist, schedds)
186 if len(schedd_dag_info) == 1:
187 # Extract the DAG info without altering the results of the query.
188 schedd_name = next(iter(schedd_dag_info))
189 dag_id = next(iter(schedd_dag_info[schedd_name]))
190 dag_ad = schedd_dag_info[schedd_name][dag_id]
192 # If the provided workflow id does not correspond to the one extracted
193 # from the DAGMan log file in the submit directory, rerun the query
194 # with the id found in the file.
195 #
196 # This is to cover the situation in which the user provided the old job
197 # id of a restarted run.
198 try:
199 path_dag_id, _ = read_dag_log(dag_ad["Iwd"])
200 except FileNotFoundError as exc:
201 # At the moment missing DAGMan log is pretty much a fatal error.
202 # So empty the DAG info to finish early (see the if statement
203 # below).
204 schedd_dag_info.clear()
205 messages.append(f"Cannot create the report for '{dag_id}': {exc}")
206 else:
207 if path_dag_id != dag_id:
208 schedd_dag_info = _get_info_from_schedd(path_dag_id, hist, schedds)
209 messages.append(
210 f"WARNING: Found newer workflow executions in same submit directory as id '{dag_id}'. "
211 "This normally occurs when a run is restarted. The report shown is for the most "
212 f"recent status with run id '{path_dag_id}'"
213 )
215 if len(schedd_dag_info) == 0:
216 run_reports = {}
217 elif len(schedd_dag_info) == 1:
218 _, dag_info = schedd_dag_info.popitem()
219 dag_id, dag_ad = dag_info.popitem()
221 # Create a mapping between jobs and their classads. The keys will
222 # be of format 'ClusterId.ProcId'.
223 job_info = {dag_id: dag_ad}
225 # Find jobs (nodes) belonging to that DAGMan job.
226 job_constraint = f"DAGManJobId == {int(float(dag_id))}"
227 schedd_job_info = condor_search(constraint=job_constraint, hist=hist, schedds=schedds)
228 if schedd_job_info:
229 _, node_info = schedd_job_info.popitem()
230 job_info.update(node_info)
232 # Collect additional pieces of information about jobs using HTCondor
233 # files in the submission directory.
234 _, path_jobs, message = _get_info_from_path(dag_ad["Iwd"])
235 _update_jobs(job_info, path_jobs)
236 if message:
237 messages.append(message)
238 run_reports = _create_detailed_report_from_jobs(dag_id, job_info)
239 else:
240 ids = [ad["GlobalJobId"] for dag_info in schedd_dag_info.values() for ad in dag_info.values()]
241 message = (
242 f"More than one job matches id '{wms_workflow_id}', "
243 f"their global ids are: {', '.join(ids)}. Rerun with one of the global ids"
244 )
245 messages.append(message)
246 run_reports = {}
248 message = "\n".join(messages)
249 return run_reports, message
252def _get_info_from_schedd(
253 wms_workflow_id: str, hist: float, schedds: dict[str, htcondor.Schedd]
254) -> dict[str, dict[str, dict[str, Any]]]:
255 """Gather run information from HTCondor.
257 Parameters
258 ----------
259 wms_workflow_id : `str`
260 Limit to specific run based on id.
261 hist : `float`
262 Limit history search to this many days.
263 schedds : `dict` [ `str`, `htcondor.Schedd` ]
264 HTCondor schedulers which to query for job information. If empty
265 dictionary, all queries will be run against the local scheduler only.
267 Returns
268 -------
269 schedd_dag_info : `dict` [`str`, `dict` [`str`, `dict` [`str` Any]]]
270 Information about jobs satisfying the search criteria where for each
271 Scheduler, local HTCondor job ids are mapped to their respective
272 classads.
273 """
274 _LOG.debug("_get_info_from_schedd: id=%s, hist=%s, schedds=%s", wms_workflow_id, hist, schedds)
276 dag_constraint = 'regexp("dagman$", Cmd)'
277 try:
278 cluster_id = int(float(wms_workflow_id))
279 except ValueError:
280 dag_constraint += f' && GlobalJobId == "{wms_workflow_id}"'
281 else:
282 dag_constraint += f" && ClusterId == {cluster_id}"
284 # With the current implementation of the condor_* functions the query
285 # will always return only one match per Scheduler.
286 #
287 # Even in the highly unlikely situation where HTCondor history (which
288 # condor_search queries too) is long enough to have jobs from before
289 # the cluster ids were rolled over (and as a result there is more then
290 # one job with the same cluster id) they will not show up in
291 # the results.
292 schedd_dag_info = condor_search(constraint=dag_constraint, hist=hist, schedds=schedds)
293 return schedd_dag_info
296def _get_info_from_path(wms_path: str | os.PathLike) -> tuple[str, dict[str, dict[str, Any]], str]:
297 """Gather run information from a given run directory.
299 Parameters
300 ----------
301 wms_path : `str` or `os.PathLike`
302 Directory containing HTCondor files.
304 Returns
305 -------
306 wms_workflow_id : `str`
307 The run id which is a DAGman job id.
308 jobs : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
309 Information about jobs read from files in the given directory.
310 The key is the HTCondor id and the value is a dictionary of HTCondor
311 keys and values.
312 message : `str`
313 Message to be printed with the summary report.
314 """
315 # Ensure path is absolute, in particular for folks helping
316 # debug failures that need to dig around submit files.
317 wms_path = Path(wms_path).resolve()
319 messages = []
320 try:
321 wms_workflow_id, jobs = read_dag_log(wms_path)
322 _LOG.debug("_get_info_from_path: from dag log %s = %s", wms_workflow_id, jobs)
323 _update_jobs(jobs, read_node_status(wms_path))
324 _LOG.debug("_get_info_from_path: after node status %s = %s", wms_workflow_id, jobs)
326 # Add more info for DAGman job
327 job = jobs[wms_workflow_id]
328 job.update(read_dag_status(wms_path))
330 job["total_jobs"], job["state_counts"] = _get_state_counts_from_jobs(wms_workflow_id, jobs)
331 if "bps_run" not in job: 331 ↛ 334line 331 didn't jump to line 334 because the condition on line 331 was always true
332 _add_run_info(wms_path, job)
334 message = htc_check_dagman_output(wms_path)
335 if message: 335 ↛ 336line 335 didn't jump to line 336 because the condition on line 335 was never true
336 messages.append(message)
337 _LOG.debug(
338 "_get_info: id = %s, total_jobs = %s", wms_workflow_id, jobs[wms_workflow_id]["total_jobs"]
339 )
341 # Add extra pieces of information which cannot be found in HTCondor
342 # generated files like 'GlobalJobId'.
343 #
344 # Do not treat absence of this file as a serious error. Neither runs
345 # submitted with earlier versions of the plugin nor the runs submitted
346 # with Pegasus plugin will have it at the moment. However, once enough
347 # time passes and Pegasus plugin will have its own report() method
348 # (instead of sneakily using HTCondor's one), the lack of that file
349 # should be treated as seriously as lack of any other file.
350 try:
351 _, job_info = read_dag_info(wms_path)
352 except FileNotFoundError as exc:
353 message = f"Warn: Some information may not be available: {exc}"
354 messages.append(message)
355 else:
356 schedd_name = next(iter(job_info))
357 job_ad = next(iter(job_info[schedd_name].values()))
358 job.update(job_ad)
359 except FileNotFoundError as err:
360 message = f"Could not find HTCondor files in '{wms_path}' ({err})"
361 _LOG.debug(message)
362 messages.append(message)
363 try:
364 message = htc_check_dagman_output(wms_path)
365 except FileNotFoundError as err:
366 message = (
367 f"Could not find DAGMan standard output file in '{wms_path}'.\n"
368 "Check that path is a valid HTCondor submission directory.\n"
369 "Check that the condor_dagman executable path is accessible from the AP machine."
370 )
371 if message:
372 messages.append(message)
373 wms_workflow_id = MISSING_ID
374 jobs = {}
376 # Add more condor_q-like info.
377 for job in jobs.values():
378 htc_tweak_log_info(wms_path, job)
380 message = "\n".join([msg for msg in messages if msg])
381 _LOG.debug("wms_workflow_id = %s, jobs = %s", wms_workflow_id, jobs.keys())
382 _LOG.debug("message = %s", message)
383 return wms_workflow_id, jobs, message
386def _create_detailed_report_from_jobs(
387 wms_workflow_id: str, jobs: dict[str, dict[str, Any]]
388) -> dict[str, WmsRunReport]:
389 """Gather run information to be used in generating summary reports.
391 Parameters
392 ----------
393 wms_workflow_id : `str`
394 The run id to create the report for.
395 jobs : `dict` [`str`, `dict` [`str`, Any]]
396 Mapping HTCondor job id to job information.
398 Returns
399 -------
400 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
401 Run information for the detailed report. The key is the given HTCondor
402 id and the value is a collection of report information for that run.
403 """
404 _LOG.debug("_create_detailed_report: id = %s, job = %s", wms_workflow_id, jobs[wms_workflow_id])
406 dag_ad = jobs[wms_workflow_id]
408 report = WmsRunReport(
409 wms_id=f"{dag_ad['ClusterId']}.{dag_ad['ProcId']}",
410 global_wms_id=dag_ad.get("GlobalJobId", "MISS"),
411 path=dag_ad["Iwd"],
412 label=dag_ad.get("bps_job_label", "MISS"),
413 run=dag_ad.get("bps_run", "MISS"),
414 site=dag_ad.get("bps_runsite", ""),
415 project=dag_ad.get("bps_project", "MISS"),
416 campaign=dag_ad.get("bps_campaign", "MISS"),
417 payload=dag_ad.get("bps_payload", "MISS"),
418 operator=_get_owner(dag_ad),
419 run_summary=_get_run_summary(dag_ad),
420 state=_htc_status_to_wms_state(dag_ad),
421 total_number_jobs=0,
422 jobs=[],
423 job_state_counts=dict.fromkeys(WmsStates, 0),
424 exit_code_summary={},
425 )
427 payload_jobs = {} # keep track for later processing
428 specific_info = WmsSpecificInfo()
429 for job_id, job_ad in jobs.items():
430 if job_ad.get("wms_node_type", WmsNodeType.UNKNOWN) in [WmsNodeType.PAYLOAD, WmsNodeType.FINAL]:
431 try:
432 name = job_ad.get("DAGNodeName", job_id)
433 wms_state = _htc_status_to_wms_state(job_ad)
434 job_report = WmsJobReport(
435 wms_id=job_id,
436 name=name,
437 label=job_ad.get("bps_job_label", pegasus_name_to_label(name)),
438 state=wms_state,
439 )
440 if job_report.label == "init": 440 ↛ 441line 440 didn't jump to line 441 because the condition on line 440 was never true
441 job_report.label = "pipetaskInit"
442 report.job_state_counts[wms_state] += 1
443 report.jobs.append(job_report)
444 payload_jobs[job_id] = job_ad
445 except KeyError as ex:
446 _LOG.error("Job missing key '%s': %s", str(ex), job_ad)
447 raise
448 elif is_service_job(job_ad):
449 _LOG.debug(
450 "Found service job: id='%s', name='%s', label='%s', NodeStatus='%s', JobStatus='%s'",
451 job_id,
452 job_ad["DAGNodeName"],
453 job_ad.get("bps_job_label", "MISS"),
454 job_ad.get("NodeStatus", "MISS"),
455 job_ad.get("JobStatus", "MISS"),
456 )
457 _add_service_job_specific_info(job_ad, specific_info)
459 report.total_number_jobs = len(payload_jobs)
460 report.exit_code_summary = _get_exit_code_summary(payload_jobs)
461 if specific_info:
462 report.specific_info = specific_info
464 # Workflow will exit with non-zero DAG_STATUS if problem with
465 # any of the wms jobs. So change FAILED to SUCCEEDED if all
466 # payload jobs SUCCEEDED.
467 if report.total_number_jobs == report.job_state_counts[WmsStates.SUCCEEDED]:
468 report.state = WmsStates.SUCCEEDED
470 run_reports = {report.wms_id: report}
471 _LOG.debug("_create_detailed_report: run_reports = %s", run_reports)
472 return run_reports
475def _add_service_job_specific_info(job_ad: dict[str, Any], specific_info: WmsSpecificInfo) -> None:
476 """Generate report information for service job.
478 Parameters
479 ----------
480 job_ad : `dict` [`str`, `~typing.Any`]
481 Provisioning job information.
482 specific_info : `lsst.ctrl.bps.WmsSpecificInfo`
483 Where to add message.
484 """
485 status_details = ""
486 job_status = _htc_status_to_wms_state(job_ad)
488 # Service jobs in queue are deleted when DAG is done.
489 # To get accurate status, need to check other info.
490 if (
491 job_status == WmsStates.DELETED
492 and "Reason" in job_ad
493 and (
494 "Removed by DAGMan" in job_ad["Reason"]
495 or "removed because <OtherJobRemoveRequirements = DAGManJobId =?=" in job_ad["Reason"]
496 or "DAG is exiting and writing rescue file." in job_ad["Reason"]
497 )
498 ):
499 if "HoldReason" in job_ad:
500 # HoldReason exists even if released, so check.
501 if "job_released_time" in job_ad and job_ad["job_held_time"] < job_ad["job_released_time"]:
502 # If released, assume running until deleted.
503 job_status = WmsStates.SUCCEEDED
504 status_details = ""
505 else:
506 # If job held when deleted by DAGMan, still want to
507 # report hold reason
508 status_details = f"(Job was held for the following reason: {job_ad['HoldReason']})"
510 else:
511 job_status = WmsStates.SUCCEEDED
512 elif job_status == WmsStates.SUCCEEDED:
513 status_details = "(Note: Finished before workflow.)"
514 elif job_status == WmsStates.HELD:
515 status_details = f"({job_ad['HoldReason']})"
517 template = "Status of {job_name}: {status} {status_details}"
518 context = {
519 "job_name": job_ad["DAGNodeName"],
520 "status": job_status.name,
521 "status_details": status_details,
522 }
523 specific_info.add_message(template=template, context=context)
526def _summary_report(user, hist, pass_thru, schedds=None):
527 """Gather run information to be used in generating summary reports.
529 Parameters
530 ----------
531 user : `str`
532 Run lookup restricted to given user.
533 hist : `float`
534 How many previous days to search for run information.
535 pass_thru : `str`
536 Advanced users can define the HTCondor constraint to be used
537 when searching queue and history.
539 Returns
540 -------
541 run_reports : `dict` [`str`, `lsst.ctrl.bps.WmsRunReport`]
542 Run information for the summary report. The keys are HTCondor ids and
543 the values are collections of report information for each run.
544 message : `str`
545 Message to be printed with the summary report.
546 """
547 # only doing summary report so only look for dagman jobs
548 if pass_thru:
549 constraint = pass_thru
550 else:
551 # Notes:
552 # * bps_isjob == 'True' isn't getting set for DAG jobs that are
553 # manually restarted.
554 # * Any job with DAGManJobID isn't a DAG job
555 constraint = 'bps_isjob == "True" && JobUniverse == 7 && DAGManJobID =?= Undefined'
556 if user:
557 constraint += f' && (Owner == "{user}" || bps_operator == "{user}")'
559 job_info = condor_search(constraint=constraint, hist=hist, schedds=schedds)
561 # Have list of DAGMan jobs, need to get run_report info.
562 run_reports = {}
563 msg = ""
564 for jobs in job_info.values():
565 for job_id, job in jobs.items():
566 total_jobs, state_counts = _get_state_counts_from_dag_job(job)
567 # If didn't get from queue information (e.g., Kerberos bug),
568 # try reading from file.
569 if total_jobs == 0:
570 try:
571 job.update(read_dag_status(job["Iwd"]))
572 total_jobs, state_counts = _get_state_counts_from_dag_job(job)
573 except (StopIteration, FileNotFoundError):
574 pass # don't kill report can't find htcondor files
576 if "bps_run" not in job:
577 _add_run_info(job["Iwd"], job)
578 report = WmsRunReport(
579 wms_id=job_id,
580 global_wms_id=job["GlobalJobId"],
581 path=job["Iwd"],
582 label=job.get("bps_job_label", "MISS"),
583 run=job.get("bps_run", "MISS"),
584 project=job.get("bps_project", "MISS"),
585 campaign=job.get("bps_campaign", "MISS"),
586 payload=job.get("bps_payload", "MISS"),
587 site=job.get("bps_runsite", ""),
588 operator=_get_owner(job),
589 run_summary=_get_run_summary(job),
590 state=_htc_status_to_wms_state(job),
591 jobs=[],
592 total_number_jobs=total_jobs,
593 job_state_counts=state_counts,
594 )
595 run_reports[report.global_wms_id] = report
597 return run_reports, msg
600def _add_run_info(wms_path, job):
601 """Find BPS run information elsewhere for runs without bps attributes.
603 Parameters
604 ----------
605 wms_path : `str`
606 Path to submit files for the run.
607 job : `dict` [`str`, `~typing.Any`]
608 HTCondor dag job information.
610 Raises
611 ------
612 StopIteration
613 If cannot find file it is looking for. Permission errors are
614 caught and job's run is marked with error.
615 """
616 path = Path(wms_path) / "jobs"
617 try:
618 subfile = next(path.glob("**/*.sub"))
619 except (StopIteration, PermissionError):
620 job["bps_run"] = "Unavailable"
621 else:
622 _LOG.debug("_add_run_info: subfile = %s", subfile)
623 try:
624 with open(subfile, encoding="utf-8") as fh:
625 for line in fh:
626 if line.startswith("+bps_"):
627 m = re.match(r"\+(bps_[^\s]+)\s*=\s*(.+)$", line)
628 if m:
629 _LOG.debug("Matching line: %s", line)
630 job[m.group(1)] = m.group(2).replace('"', "")
631 else:
632 _LOG.debug("Could not parse attribute: %s", line)
633 except PermissionError:
634 job["bps_run"] = "PermissionError"
635 _LOG.debug("After adding job = %s", job)
638def _get_owner(job):
639 """Get the owner of a dag job.
641 Parameters
642 ----------
643 job : `dict` [`str`, `~typing.Any`]
644 HTCondor dag job information.
646 Returns
647 -------
648 owner : `str`
649 Owner of the dag job.
650 """
651 owner = job.get("bps_operator", None)
652 if not owner: 652 ↛ 653line 652 didn't jump to line 653 because the condition on line 652 was never true
653 owner = job.get("Owner", None)
654 if not owner:
655 _LOG.warning("Could not get Owner from htcondor job: %s", job)
656 owner = "MISS"
657 return owner
660def _get_run_summary(job):
661 """Get the run summary for a job.
663 Parameters
664 ----------
665 job : `dict` [`str`, `~typing.Any`]
666 HTCondor dag job information.
668 Returns
669 -------
670 summary : `str`
671 Number of jobs per PipelineTask label in approximate pipeline order.
672 Format: <label>:<count>[;<label>:<count>]+
673 """
674 summary = job.get("bps_job_summary", job.get("bps_run_summary", None))
675 if not summary:
676 summary, _, _ = summarize_dag(job["Iwd"])
677 if not summary:
678 _LOG.warning("Could not get run summary for htcondor job: %s", job)
679 _LOG.debug("_get_run_summary: summary=%s", summary)
681 # Workaround sometimes using init vs pipetaskInit
682 summary = summary.replace("init:", "pipetaskInit:")
684 if "pegasus_version" in job and "pegasus" not in summary: 684 ↛ 685line 684 didn't jump to line 685 because the condition on line 684 was never true
685 summary += ";pegasus:0"
687 return summary
690def _get_exit_code_summary(jobs):
691 """Get the exit code summary for a run.
693 Parameters
694 ----------
695 jobs : `dict` [`str`, `dict` [`str`, Any]]
696 Mapping HTCondor job id to job information.
698 Returns
699 -------
700 summary : `dict` [`str`, `list` [`int`]]
701 Jobs' exit codes per job label.
702 """
703 summary = {}
704 for job_id, job_ad in jobs.items():
705 job_label = job_ad["bps_job_label"]
706 summary.setdefault(job_label, [])
707 try:
708 exit_code = 0
709 job_status = job_ad["JobStatus"]
710 match job_status:
711 case htcondor.JobStatus.COMPLETED | htcondor.JobStatus.HELD | htcondor.JobStatus.REMOVED:
712 exit_code = job_ad["ExitSignal"] if job_ad["ExitBySignal"] else job_ad["ExitCode"]
713 case (
714 htcondor.JobStatus.IDLE
715 | htcondor.JobStatus.RUNNING
716 | htcondor.JobStatus.TRANSFERRING_OUTPUT
717 | htcondor.JobStatus.SUSPENDED
718 ):
719 pass
720 case _:
721 _LOG.debug("Unknown 'JobStatus' value ('%d') in classad for job '%s'", job_status, job_id)
722 if exit_code != 0:
723 summary[job_label].append(exit_code)
724 except KeyError as ex:
725 _LOG.debug("Attribute '%s' not found in the classad for job '%s'", ex, job_id)
726 return summary
729def _get_state_counts_from_jobs(
730 wms_workflow_id: str, jobs: dict[str, dict[str, Any]]
731) -> tuple[int, dict[WmsStates, int]]:
732 """Count number of jobs per WMS state.
734 The workflow job and the service jobs are excluded from the count.
736 Parameters
737 ----------
738 wms_workflow_id : `str`
739 HTCondor job id.
740 jobs : `dict [`dict` [`str`, `~typing.Any`]]
741 HTCondor dag job information.
743 Returns
744 -------
745 total_count : `int`
746 Total number of dag nodes.
747 state_counts : `dict` [`lsst.ctrl.bps.WmsStates`, `int`]
748 Keys are the different WMS states and values are counts of jobs
749 that are in that WMS state.
750 """
751 state_counts = dict.fromkeys(WmsStates, 0)
752 for job_id, job_ad in jobs.items():
753 if job_id != wms_workflow_id and job_ad.get("wms_node_type", WmsNodeType.UNKNOWN) in [
754 WmsNodeType.PAYLOAD,
755 WmsNodeType.FINAL,
756 ]:
757 state_counts[_htc_status_to_wms_state(job_ad)] += 1
758 total_count = sum(state_counts.values())
760 return total_count, state_counts
763def _get_state_counts_from_dag_job(job):
764 """Count number of jobs per WMS state.
766 Parameters
767 ----------
768 job : `dict` [`str`, `~typing.Any`]
769 HTCondor dag job information.
771 Returns
772 -------
773 total_count : `int`
774 Total number of dag nodes.
775 state_counts : `dict` [`lsst.ctrl.bps.WmsStates`, `int`]
776 Keys are the different WMS states and values are counts of jobs
777 that are in that WMS state.
778 """
779 _LOG.debug("_get_state_counts_from_dag_job: job = %s %s", type(job), len(job))
780 state_counts = dict.fromkeys(WmsStates, 0)
781 if "DAG_NodesReady" in job: 781 ↛ 793line 781 didn't jump to line 793 because the condition on line 781 was always true
782 state_counts = {
783 WmsStates.UNREADY: job.get("DAG_NodesUnready", 0),
784 WmsStates.READY: job.get("DAG_NodesReady", 0),
785 WmsStates.HELD: job.get("DAG_JobsHeld", 0),
786 WmsStates.SUCCEEDED: job.get("DAG_NodesDone", 0),
787 WmsStates.FAILED: job.get("DAG_NodesFailed", 0),
788 WmsStates.PRUNED: job.get("DAG_NodesFutile", 0),
789 WmsStates.MISFIT: job.get("DAG_NodesPre", 0) + job.get("DAG_NodesPost", 0),
790 }
791 total_jobs = job.get("DAG_NodesTotal")
792 _LOG.debug("_get_state_counts_from_dag_job: from DAG_* keys, total_jobs = %s", total_jobs)
793 elif "NodesFailed" in job:
794 state_counts = {
795 WmsStates.UNREADY: job.get("NodesUnready", 0),
796 WmsStates.READY: job.get("NodesReady", 0),
797 WmsStates.HELD: job.get("JobProcsHeld", 0),
798 WmsStates.SUCCEEDED: job.get("NodesDone", 0),
799 WmsStates.FAILED: job.get("NodesFailed", 0),
800 WmsStates.PRUNED: job.get("NodesFutile", 0),
801 WmsStates.MISFIT: job.get("NodesPre", 0) + job.get("NodesPost", 0),
802 }
803 try:
804 total_jobs = job.get("NodesTotal")
805 except KeyError as ex:
806 _LOG.error("Job missing %s. job = %s", str(ex), job)
807 raise
808 _LOG.debug("_get_state_counts_from_dag_job: from NODES* keys, total_jobs = %s", total_jobs)
809 else:
810 # With Kerberos job auth and Kerberos bug, if warning would be printed
811 # for every DAG.
812 _LOG.debug("Can't get job state counts %s", job["Iwd"])
813 total_jobs = 0
815 _LOG.debug("total_jobs = %s, state_counts: %s", total_jobs, state_counts)
816 return total_jobs, state_counts
819def _update_jobs(jobs1, jobs2):
820 """Update jobs1 with info in jobs2.
822 (Basically an update for nested dictionaries.)
824 Parameters
825 ----------
826 jobs1 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
827 HTCondor job information to be updated.
828 jobs2 : `dict` [`str`, `dict` [`str`, `~typing.Any`]]
829 Additional HTCondor job information.
830 """
831 for job_id, job_ad in jobs2.items():
832 if job_id in jobs1: 832 ↛ 833line 832 didn't jump to line 833 because the condition on line 832 was never true
833 jobs1[job_id].update(job_ad)
834 else:
835 jobs1[job_id] = job_ad
838def is_service_job(job_ad: dict[str, Any]) -> bool:
839 """Determine if a job is a service one.
841 Parameters
842 ----------
843 job_ad : `dict` [`str`, Any]
844 Information about an HTCondor job.
846 Returns
847 -------
848 is_service_job : `bool`
849 True if the job is a service one, false otherwise.
851 Notes
852 -----
853 At the moment, HTCondor does not provide a native way to distinguish
854 between payload and service jobs in the workflow. This code depends
855 on read_node_status adding wms_node_type.
856 """
857 return job_ad.get("wms_node_type", WmsNodeType.UNKNOWN) == WmsNodeType.SERVICE