Coverage for python/lsst/ctrl/bps/bps_reports.py: 92%

186 statements  

« prev     ^ index     » next       coverage.py v7.16.2, created at 2026-09-28 09:22 +0000

1# This file is part of ctrl_bps. 

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"""Classes and functions used in reporting run status.""" 

29 

30__all__ = [ 

31 "DISPLAY_SUMMARY_FIELDS", 

32 "BaseRunReport", 

33 "DetailedRunReport", 

34 "ExitCodesReport", 

35 "SummaryRunReport", 

36 "compile_code_summary", 

37 "compile_job_summary", 

38] 

39 

40import abc 

41import logging 

42 

43from astropy.table import Table 

44 

45from .wms_service import WmsRunReport, WmsStates 

46 

47_LOG = logging.getLogger(__name__) 

48 

49 

50DISPLAY_SUMMARY_FIELDS = [ 

51 ("X", "S"), 

52 ("STATE", "S"), 

53 ("%S", "S"), 

54 ("ID", "S"), 

55 ("OPERATOR", "S"), 

56 ("PROJECT", "S"), 

57 ("CAMPAIGN", "S"), 

58 ("SITE", "S"), 

59 ("PAYLOAD", "S"), 

60 ("RUN", "S"), 

61] 

62 

63 

64class BaseRunReport(abc.ABC): 

65 """The base class representing a run report. 

66 

67 Parameters 

68 ---------- 

69 fields : `list` [ `tuple` [ `str`, `str`]] 

70 The list of column specification, fields, to include in the report. 

71 Each field has a name and a type. 

72 """ 

73 

74 def __init__(self, fields): 

75 self._table = Table(dtype=fields) 

76 self._msg = None 

77 

78 def __eq__(self, other): 

79 if isinstance(other, BaseRunReport): 79 ↛ 81line 79 didn't jump to line 81 because the condition on line 79 was always true

80 return self._table.pformat() == other._table.pformat() 

81 return False 

82 

83 def __len__(self): 

84 """Return the number of runs in the report.""" 

85 return len(self._table) 

86 

87 def __str__(self): 

88 lines = list(self._table.pformat(max_lines=-1, max_width=-1)) 

89 return "\n".join(lines) 

90 

91 @property 

92 def message(self): 

93 """Extra information a method need to pass to its caller (`str`).""" 

94 return self._msg 

95 

96 def clear(self): 

97 """Remove all entries from the report.""" 

98 self._msg = None 

99 self._table.remove_rows(slice(len(self))) 

100 

101 def sort(self, columns, ascending=True): 

102 """Sort the report entries according to one or more keys. 

103 

104 Parameters 

105 ---------- 

106 columns : `str` | `list` [ `str` ] 

107 The column(s) to order the report by. 

108 ascending : `bool`, optional 

109 Sort report entries in ascending order, default. 

110 

111 Raises 

112 ------ 

113 AttributeError 

114 Raised if supplied with non-existent column(s). 

115 """ 

116 if isinstance(columns, str): 116 ↛ 118line 116 didn't jump to line 118 because the condition on line 116 was always true

117 columns = [columns] 

118 unknown_keys = set(columns) - set(self._table.colnames) 

119 if unknown_keys: 

120 raise AttributeError( 

121 f"cannot sort the report entries: column(s) {', '.join(unknown_keys)} not found" 

122 ) 

123 self._table.sort(keys=columns, reverse=not ascending) 

124 

125 @classmethod 

126 def from_table(cls, table): 

127 """Create a report from a table. 

128 

129 Parameters 

130 ---------- 

131 table : `astropy.table.Table` 

132 Information about a run in a tabular form. 

133 

134 Returns 

135 ------- 

136 inst : `lsst.ctrl.bps.bps_reports.BaseRunReport` 

137 A report created based on the information in the provided table. 

138 """ 

139 inst = cls(table.dtype.descr) 

140 inst._table = table.copy() 

141 return inst 

142 

143 @abc.abstractmethod 

144 def add(self, run_report, use_global_id=False): 

145 """Add a single run info to the report. 

146 

147 Parameters 

148 ---------- 

149 run_report : `lsst.ctrl.bps.WmsRunReport` 

150 Information for single run. 

151 use_global_id : `bool`, optional 

152 If set, use global run id. Defaults to False which means that 

153 the local id will be used instead. 

154 

155 Only applicable in the context of a WMS using distributed job 

156 queues (e.g., HTCondor). 

157 """ 

158 

159 

160class SummaryRunReport(BaseRunReport): 

161 """A summary run report.""" 

162 

163 def add(self, run_report, use_global_id=False): 

164 # Docstring inherited from the base class. 

165 

166 # Flag any running workflow that might need human attention. 

167 run_flag = " " 

168 if run_report.state == WmsStates.RUNNING: 

169 if run_report.job_state_counts.get(WmsStates.HELD, 0): 

170 run_flag = "H" 

171 elif run_report.job_state_counts.get(WmsStates.DELETED, 0): 

172 run_flag = "D" 

173 elif run_report.job_state_counts.get(WmsStates.FAILED, 0): 

174 run_flag = "F" 

175 

176 # Estimate success rate. 

177 percent_succeeded = "UNK" 

178 _LOG.debug("total_number_jobs = %s", run_report.total_number_jobs) 

179 _LOG.debug("run_report.job_state_counts = %s", run_report.job_state_counts) 

180 if run_report.total_number_jobs: 180 ↛ 185line 180 didn't jump to line 185 because the condition on line 180 was always true

181 succeeded = run_report.job_state_counts.get(WmsStates.SUCCEEDED, 0) 

182 _LOG.debug("succeeded = %s", succeeded) 

183 percent_succeeded = f"{int(succeeded / run_report.total_number_jobs * 100)}" 

184 

185 row = ( 

186 run_flag, 

187 run_report.state.name, 

188 percent_succeeded, 

189 run_report.global_wms_id if use_global_id else run_report.wms_id, 

190 run_report.operator if run_report.operator is not None else "", 

191 run_report.project if run_report.project is not None else "", 

192 run_report.campaign if run_report.campaign is not None else "", 

193 run_report.site if run_report.site is not None else "", 

194 run_report.payload if run_report.payload is not None else "", 

195 run_report.run, 

196 ) 

197 try: 

198 self._table.add_row(row) 

199 except ValueError as ex: 

200 _LOG.error("Error when adding summary report row: %s (%s)", repr(ex), row) 

201 

202 

203class DetailedRunReport(BaseRunReport): 

204 """A detailed run report.""" 

205 

206 def add(self, run_report, use_global_id=False): 

207 # Docstring inherited from the base class. 

208 

209 # If run summary exists, use it to get the reference job counts. 

210 by_label_expected = {} 

211 if run_report.run_summary: 

212 for part in run_report.run_summary.split(";"): 

213 label, count = part.split(":") 

214 by_label_expected[label] = int(count) 

215 

216 total = ["TOTAL"] 

217 total.extend([run_report.job_state_counts[state] for state in WmsStates]) 

218 total.append(sum(by_label_expected.values()) if by_label_expected else run_report.total_number_jobs) 

219 self._table.add_row(total) 

220 

221 job_summary = run_report.job_summary 

222 if job_summary is None: 

223 id_ = run_report.global_wms_id if use_global_id else run_report.wms_id 

224 self._msg = f"WARNING: Job summary for run '{id_}' not available, report may be incomplete." 

225 return 

226 

227 if by_label_expected: 

228 job_order = list(by_label_expected) 

229 else: 

230 job_order = sorted(job_summary) 

231 self._msg = "WARNING: Could not determine order of pipeline, instead sorted alphabetically." 

232 for label in job_order: 

233 try: 

234 counts = job_summary[label] 

235 except KeyError: 

236 counts = dict.fromkeys(WmsStates, -1) 

237 else: 

238 if label in by_label_expected: 

239 already_counted = sum(counts.values()) 

240 if already_counted != by_label_expected[label]: 240 ↛ 241line 240 didn't jump to line 241 because the condition on line 240 was never true

241 counts[WmsStates.UNREADY] += by_label_expected[label] - already_counted 

242 

243 run = [label] 

244 run.extend([counts[state] for state in WmsStates]) 

245 run.append(by_label_expected[label] if by_label_expected else -1) 

246 self._table.add_row(run) 

247 

248 def __str__(self): 

249 alignments = ["<"] + [">"] * (len(self._table.colnames) - 1) 

250 lines = list(self._table.pformat(max_lines=-1, max_width=-1, align=alignments)) 

251 lines.insert(3, lines[1]) 

252 return str("\n".join(lines)) 

253 

254 

255class ExitCodesReport(BaseRunReport): 

256 """An extension of run report to give information about 

257 error handling from the wms service. 

258 """ 

259 

260 def add(self, run_report: WmsRunReport, use_global_id: bool = False) -> None: 

261 # Docstring inherited from the base class. 

262 

263 exit_code_summary = run_report.exit_code_summary 

264 if not exit_code_summary: 

265 id_ = run_report.global_wms_id if use_global_id else run_report.wms_id 

266 self._msg = f"WARNING: Exit code summary for run '{id_}' not available, report may be incomplete." 

267 return 

268 

269 warnings = [] 

270 

271 # If available, use label ordering from the run summary as it should 

272 # reflect the ordering of the pipetasks in the pipeline. 

273 labels = [] 

274 if run_report.run_summary: 

275 for part in run_report.run_summary.split(";"): 

276 label, _ = part.split(":") 

277 labels.append(label) 

278 if not labels: 

279 labels = sorted(exit_code_summary) 

280 warnings.append("WARNING: Could not determine order of pipeline, instead sorted alphabetically.") 

281 

282 # Payload (e.g. pipetask) error codes: 

283 # * 1: general failure, 

284 # * 2: command line error (e.g. unknown command and/or option). 

285 pyld_error_codes = {1, 2} 

286 

287 missing_labels = set() 

288 for label in labels: 

289 try: 

290 exit_codes = exit_code_summary[label] 

291 except KeyError: 

292 missing_labels.add(label) 

293 else: 

294 pyld_errors = [code for code in exit_codes if code in pyld_error_codes] 

295 pyld_error_count = len(pyld_errors) 

296 pyld_error_summary = ( 

297 ", ".join(sorted(str(code) for code in set(pyld_errors))) if pyld_errors else "None" 

298 ) 

299 

300 infra_errors = [code for code in exit_codes if code not in pyld_error_codes] 

301 infra_error_count = len(infra_errors) 

302 infra_error_summary = ( 

303 ", ".join(sorted(str(code) for code in set(infra_errors))) if infra_errors else "None" 

304 ) 

305 

306 run = [label, pyld_error_count, pyld_error_summary, infra_error_count, infra_error_summary] 

307 self._table.add_row(run) 

308 if missing_labels: 308 ↛ 309line 308 didn't jump to line 309 because the condition on line 308 was never true

309 warnings.append( 

310 f"WARNING: Exit code summary was not available for job labels: {', '.join(missing_labels)}" 

311 ) 

312 if warnings: 

313 self._msg = "\n".join(warnings) 

314 

315 def __str__(self): 

316 alignments = ["<"] + [">"] * (len(self._table.colnames) - 1) 

317 lines = list(self._table.pformat(max_lines=-1, max_width=-1, align=alignments)) 

318 return str("\n".join(lines)) 

319 

320 

321def compile_job_summary(report: WmsRunReport) -> list[str]: 

322 """Add a job summary to the run report if necessary. 

323 

324 If the job summary is not provided, the function will attempt to compile 

325 it from information available for individual jobs (if any) and add it to 

326 the report. If the report already includes a job summary, the function is 

327 effectively a no-op. 

328 

329 Parameters 

330 ---------- 

331 report : `lsst.ctrl.bps.WmsRunReport` 

332 Information about a single run. 

333 

334 Returns 

335 ------- 

336 warnings : `list` [`str`] 

337 List of messages describing any non-critical issues encountered during 

338 processing. Empty if none. 

339 """ 

340 warnings: list[str] = [] 

341 

342 # If the job summary already exists, exit early. 

343 if report.job_summary: 

344 return warnings 

345 

346 if report.jobs: 

347 job_summary = {} 

348 by_label = group_jobs_by_label(report.jobs) 

349 for label, job_group in by_label.items(): 

350 by_label_state = group_jobs_by_state(job_group) 

351 _LOG.debug("by_label_state = %s", by_label_state) 

352 counts = {state: len(jobs) for state, jobs in by_label_state.items()} 

353 job_summary[label] = counts 

354 report.job_summary = job_summary 

355 else: 

356 warnings.append("information about individual jobs not available") 

357 

358 return warnings 

359 

360 

361def compile_code_summary(report: WmsRunReport) -> list[str]: 

362 """Add missing entries to the exit code summary if necessary. 

363 

364 A WMS plugin may exclude job labels for which there are no failures from 

365 the exit code summary. The function will attempt to use the job summary, 

366 if available, to add missing entries for these labels. 

367 

368 Parameters 

369 ---------- 

370 report : `lsst.ctrl.bps.WmsRunReport` 

371 Information about a single run. 

372 

373 Returns 

374 ------- 

375 warnings : `list` [`str`] 

376 List of messages describing any non-critical issues encountered during 

377 processing. Empty if none. 

378 """ 

379 warnings: list[str] = [] 

380 

381 # If the job summary is not available, exit early. 

382 if not report.job_summary: 

383 return warnings 

384 

385 # A shallow copy is enough here because we won't be modifying the existing 

386 # entries, only adding new ones if necessary. 

387 exit_code_summary = dict(report.exit_code_summary) if report.exit_code_summary else {} 

388 

389 # Use the job summary to add the entries for labels with no failures 

390 # *without* modifying already existing entries. 

391 failure_summary = {label: states[WmsStates.FAILED] for label, states in report.job_summary.items()} 

392 for label, count in failure_summary.items(): 

393 if count == 0: 

394 exit_code_summary.setdefault(label, []) 

395 

396 # Check if there are any discrepancies between the data in the exit code 

397 # summary and the job summary. 

398 code_summary_labels = set(exit_code_summary) 

399 failure_summary_labels = set(failure_summary) 

400 mismatches = { 

401 label 

402 for label in failure_summary_labels & code_summary_labels 

403 if len(exit_code_summary[label]) != failure_summary[label] 

404 } 

405 if mismatches: 

406 warnings.append( 

407 f"number of exit codes differs from number of failures for job labels: {', '.join(mismatches)}" 

408 ) 

409 missing = failure_summary_labels - code_summary_labels 

410 if missing: 

411 warnings.append(f"exit codes not available for job labels: {', '.join(missing)}") 

412 

413 if exit_code_summary: 413 ↛ 416line 413 didn't jump to line 416 because the condition on line 413 was always true

414 report.exit_code_summary = exit_code_summary 

415 

416 return warnings 

417 

418 

419def group_jobs_by_state(jobs): 

420 """Divide given jobs into groups based on their state value. 

421 

422 Parameters 

423 ---------- 

424 jobs : `list` [`lsst.ctrl.bps.WmsJobReport`] 

425 Jobs to divide into groups based on state. 

426 

427 Returns 

428 ------- 

429 by_state : `dict` 

430 Mapping of job state to a list of jobs. 

431 """ 

432 _LOG.debug("group_jobs_by_state: jobs=%s", jobs) 

433 by_state = {state: [] for state in WmsStates} 

434 for job in jobs: 

435 by_state[job.state].append(job) 

436 return by_state 

437 

438 

439def group_jobs_by_label(jobs): 

440 """Divide given jobs into groups based on their label value. 

441 

442 Parameters 

443 ---------- 

444 jobs : `list` [`lsst.ctrl.bps.WmsJobReport`] 

445 Jobs to divide into groups based on label. 

446 

447 Returns 

448 ------- 

449 by_label : `dict` [`str`, `list` [`lsst.ctrl.bps.WmsJobReport`]] 

450 Mapping of job state to a list of jobs. 

451 """ 

452 by_label = {} 

453 for job in jobs: 

454 group = by_label.setdefault(job.label, []) 

455 group.append(job) 

456 return by_label