Coverage for python/lsst/ctrl/bps/htcondor/report_utils.py: 61%

294 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-09-06 01:56 -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/>. 

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

125 

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

131 

132 return state, message 

133 

134 

135def _report_from_path(wms_path): 

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

137 

138 Parameters 

139 ---------- 

140 wms_path : `str` 

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

142 

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 

157 

158 

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

160 """Gather run information using workflow id. 

161 

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. 

171 

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

181 

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] 

190 

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 ) 

213 

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

219 

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} 

223 

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) 

230 

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

246 

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

248 return run_reports, message 

249 

250 

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. 

255 

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. 

265 

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) 

274 

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

282 

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 

293 

294 

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. 

297 

298 Parameters 

299 ---------- 

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

301 Directory containing HTCondor files. 

302 

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

317 

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) 

324 

325 # Add more info for DAGman job 

326 job = jobs[wms_workflow_id] 

327 job.update(read_dag_status(wms_path)) 

328 

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) 

332 

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 ) 

339 

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

374 

375 # Add more condor_q-like info. 

376 for job in jobs.values(): 

377 htc_tweak_log_info(wms_path, job) 

378 

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 

383 

384 

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. 

389 

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. 

396 

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

404 

405 dag_ad = jobs[wms_workflow_id] 

406 

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 ) 

425 

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) 

457 

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 

462 

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 

468 

469 run_reports = {report.wms_id: report} 

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

471 return run_reports 

472 

473 

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

475 """Generate report information for service job. 

476 

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) 

486 

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

508 

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

515 

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) 

523 

524 

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

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

527 

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. 

537 

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

557 

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

559 

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 

574 

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 

595 

596 return run_reports, msg 

597 

598 

599def _add_run_info(wms_path, job): 

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

601 

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. 

608 

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) 

635 

636 

637def _get_owner(job): 

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

639 

640 Parameters 

641 ---------- 

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

643 HTCondor dag job information. 

644 

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 

657 

658 

659def _get_run_summary(job): 

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

661 

662 Parameters 

663 ---------- 

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

665 HTCondor dag job information. 

666 

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) 

679 

680 # Workaround sometimes using init vs pipetaskInit 

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

682 

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" 

685 

686 return summary 

687 

688 

689def _get_exit_code_summary(jobs): 

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

691 

692 Parameters 

693 ---------- 

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

695 Mapping HTCondor job id to job information. 

696 

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 

726 

727 

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. 

732 

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

734 

735 Parameters 

736 ---------- 

737 wms_workflow_id : `str` 

738 HTCondor job id. 

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

740 HTCondor dag job information. 

741 

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

758 

759 return total_count, state_counts 

760 

761 

762def _get_state_counts_from_dag_job(job): 

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

764 

765 Parameters 

766 ---------- 

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

768 HTCondor dag job information. 

769 

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 

813 

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

815 return total_jobs, state_counts 

816 

817 

818def _update_jobs(jobs1, jobs2): 

819 """Update jobs1 with info in jobs2. 

820 

821 (Basically an update for nested dictionaries.) 

822 

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 

835 

836 

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

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

839 

840 Parameters 

841 ---------- 

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

843 Information about an HTCondor job. 

844 

845 Returns 

846 ------- 

847 is_service_job : `bool` 

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

849 

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