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-19 09:30 +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/>. 

27 

28"""Utility functions used for reporting.""" 

29 

30import logging 

31import os 

32import re 

33from pathlib import Path 

34from typing import Any 

35 

36import htcondor 

37 

38from lsst.ctrl.bps import ( 

39 WmsJobReport, 

40 WmsRunReport, 

41 WmsSpecificInfo, 

42 WmsStates, 

43) 

44 

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) 

59 

60_LOG = logging.getLogger(__name__) 

61 

62 

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. 

67 

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. 

77 

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) 

86 

87 message = "" 

88 

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 

101 

102 

103def _get_status_from_path(wms_path: str | os.PathLike) -> tuple[WmsStates, str]: 

104 """Gather run status from a given run directory. 

105 

106 Parameters 

107 ---------- 

108 wms_path : `str` | `os.PathLike` 

109 The directory containing the submit side files (e.g., HTCondor files). 

110 

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." 

126 

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]) 

132 

133 return state, message 

134 

135 

136def _report_from_path(wms_path): 

137 """Gather run information from a given run directory. 

138 

139 Parameters 

140 ---------- 

141 wms_path : `str` 

142 The directory containing the submit side files (e.g., HTCondor files). 

143 

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 

158 

159 

160def _report_from_id(wms_workflow_id, hist, schedds=None): 

161 """Gather run information using workflow id. 

162 

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. 

172 

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 = [] 

182 

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] 

191 

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 ) 

214 

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() 

220 

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} 

224 

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) 

231 

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 = {} 

247 

248 message = "\n".join(messages) 

249 return run_reports, message 

250 

251 

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. 

256 

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. 

266 

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) 

275 

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}" 

283 

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 

294 

295 

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. 

298 

299 Parameters 

300 ---------- 

301 wms_path : `str` or `os.PathLike` 

302 Directory containing HTCondor files. 

303 

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() 

318 

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) 

325 

326 # Add more info for DAGman job 

327 job = jobs[wms_workflow_id] 

328 job.update(read_dag_status(wms_path)) 

329 

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) 

333 

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 ) 

340 

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 = {} 

375 

376 # Add more condor_q-like info. 

377 for job in jobs.values(): 

378 htc_tweak_log_info(wms_path, job) 

379 

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 

384 

385 

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. 

390 

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. 

397 

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]) 

405 

406 dag_ad = jobs[wms_workflow_id] 

407 

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 ) 

426 

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) 

458 

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 

463 

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 

469 

470 run_reports = {report.wms_id: report} 

471 _LOG.debug("_create_detailed_report: run_reports = %s", run_reports) 

472 return run_reports 

473 

474 

475def _add_service_job_specific_info(job_ad: dict[str, Any], specific_info: WmsSpecificInfo) -> None: 

476 """Generate report information for service job. 

477 

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) 

487 

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']})" 

509 

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']})" 

516 

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) 

524 

525 

526def _summary_report(user, hist, pass_thru, schedds=None): 

527 """Gather run information to be used in generating summary reports. 

528 

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. 

538 

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}")' 

558 

559 job_info = condor_search(constraint=constraint, hist=hist, schedds=schedds) 

560 

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 

575 

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 

596 

597 return run_reports, msg 

598 

599 

600def _add_run_info(wms_path, job): 

601 """Find BPS run information elsewhere for runs without bps attributes. 

602 

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. 

609 

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) 

636 

637 

638def _get_owner(job): 

639 """Get the owner of a dag job. 

640 

641 Parameters 

642 ---------- 

643 job : `dict` [`str`, `~typing.Any`] 

644 HTCondor dag job information. 

645 

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 

658 

659 

660def _get_run_summary(job): 

661 """Get the run summary for a job. 

662 

663 Parameters 

664 ---------- 

665 job : `dict` [`str`, `~typing.Any`] 

666 HTCondor dag job information. 

667 

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) 

680 

681 # Workaround sometimes using init vs pipetaskInit 

682 summary = summary.replace("init:", "pipetaskInit:") 

683 

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" 

686 

687 return summary 

688 

689 

690def _get_exit_code_summary(jobs): 

691 """Get the exit code summary for a run. 

692 

693 Parameters 

694 ---------- 

695 jobs : `dict` [`str`, `dict` [`str`, Any]] 

696 Mapping HTCondor job id to job information. 

697 

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 

727 

728 

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. 

733 

734 The workflow job and the service jobs are excluded from the count. 

735 

736 Parameters 

737 ---------- 

738 wms_workflow_id : `str` 

739 HTCondor job id. 

740 jobs : `dict [`dict` [`str`, `~typing.Any`]] 

741 HTCondor dag job information. 

742 

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()) 

759 

760 return total_count, state_counts 

761 

762 

763def _get_state_counts_from_dag_job(job): 

764 """Count number of jobs per WMS state. 

765 

766 Parameters 

767 ---------- 

768 job : `dict` [`str`, `~typing.Any`] 

769 HTCondor dag job information. 

770 

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 

814 

815 _LOG.debug("total_jobs = %s, state_counts: %s", total_jobs, state_counts) 

816 return total_jobs, state_counts 

817 

818 

819def _update_jobs(jobs1, jobs2): 

820 """Update jobs1 with info in jobs2. 

821 

822 (Basically an update for nested dictionaries.) 

823 

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 

836 

837 

838def is_service_job(job_ad: dict[str, Any]) -> bool: 

839 """Determine if a job is a service one. 

840 

841 Parameters 

842 ---------- 

843 job_ad : `dict` [`str`, Any] 

844 Information about an HTCondor job. 

845 

846 Returns 

847 ------- 

848 is_service_job : `bool` 

849 True if the job is a service one, false otherwise. 

850 

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