Coverage for python/lsst/ctrl/bps/report.py: 53%

62 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-26 09:33 +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"""Supporting functions for reporting on runs submitted to a WMS. 

29 

30Note: Expectations are that future reporting effort will revolve around LSST 

31oriented database tables. 

32""" 

33 

34__all__ = ["display_report", "retrieve_report"] 

35 

36import logging 

37import sys 

38from collections.abc import Callable, Sequence 

39from typing import TextIO 

40 

41from lsst.utils import doImportType 

42 

43from .bps_reports import DISPLAY_SUMMARY_FIELDS, DetailedRunReport, ExitCodesReport, SummaryRunReport 

44from .wms_service import BaseWmsService, WmsRunReport, WmsStates 

45 

46_LOG = logging.getLogger(__name__) 

47 

48 

49def display_report( 

50 runs: list[WmsRunReport], 

51 messages: list[str], 

52 is_detailed: bool = False, 

53 is_global: bool = False, 

54 return_exit_codes: bool = False, 

55 file: TextIO = sys.stdout, 

56) -> None: 

57 """Print out summary of jobs submitted for execution. 

58 

59 Parameters 

60 ---------- 

61 runs : `list` [`str`] 

62 Runs to include in the summary. 

63 messages : `list` [`str`] 

64 Errors that happened during report and/or processing. Empty if 

65 no issues were encountered. 

66 is_detailed : `bool`, optional 

67 If set, the function prints out a detailed report including statuses 

68 of each task in the workflow grouped by task labels. By default, only 

69 a brief summary of each run is displayed. 

70 is_global : `bool`, optional 

71 If set, a global run id(s) will be used when displaying the report. 

72 By default, the report will use local run id(s). 

73 

74 Only applicable in the context of a WMS using distributed job queues 

75 (e.g., HTCondor). 

76 return_exit_codes : `bool`, optional 

77 If set, return exit codes related to jobs with a 

78 non-success status. Defaults to False, which means that only 

79 the summary state is returned. 

80 

81 Only applicable in the context of a WMS with associated 

82 handlers to return exit codes from jobs. 

83 file : TextIO 

84 File or file-like object to write the output to. 

85 """ 

86 run_brief = SummaryRunReport(DISPLAY_SUMMARY_FIELDS) 

87 

88 if is_detailed: 88 ↛ 89line 88 didn't jump to line 89 because the condition on line 88 was never true

89 fields = [(" ", "S")] + [(state.name, "i") for state in WmsStates] + [("EXPECTED", "i")] 

90 run_report = DetailedRunReport(fields) 

91 

92 for run in runs: 

93 run_brief.add(run, use_global_id=is_global) 

94 

95 run_report.add(run, use_global_id=is_global) 

96 if run_report.message: 

97 messages.append(run_report.message) 

98 

99 print(run_brief, file=file) 

100 print("\n", file=file) 

101 print(f"Path: {run.path}", file=file) 

102 print(f"Global job id: {run.global_wms_id}", file=file) 

103 if run.specific_info: 

104 print(run.specific_info, file=file) 

105 print("\n", file=file) 

106 print(run_report, file=file) 

107 

108 if return_exit_codes: 

109 fields = [ 

110 (" ", "S"), 

111 ("PAYLOAD ERROR COUNT", "i"), 

112 ("PAYLOAD ERROR CODES", "S"), 

113 ("INFRASTRUCTURE ERROR COUNT", "i"), 

114 ("INFRASTRUCTURE ERROR CODES", "S"), 

115 ] 

116 run_exits_report = ExitCodesReport(fields) 

117 run_exits_report.add(run, use_global_id=is_global) 

118 if run_exits_report.message: 

119 messages.append(run_exits_report.message) 

120 print("\n", file=file) 

121 print(run_exits_report, file=file) 

122 run_exits_report.clear() 

123 

124 run_brief.clear() 

125 run_report.clear() 

126 else: 

127 for run in runs: 

128 run_brief.add(run, use_global_id=is_global) 

129 run_brief.sort("ID") 

130 print(run_brief, file=file) 

131 

132 if messages: 132 ↛ 133line 132 didn't jump to line 133 because the condition on line 132 was never true

133 uniques = list(dict.fromkeys(messages)) 

134 print("\n".join(uniques), file=file) 

135 print("\n", file=file) 

136 

137 

138def retrieve_report( 

139 wms_service_fqn: str, 

140 *, 

141 run_id: str | None = None, 

142 user: str | None = None, 

143 hist: float | None = None, 

144 pass_thru: str | None = None, 

145 is_global: bool = False, 

146 return_exit_codes: bool = False, 

147 postprocessors: Sequence[Callable[[WmsRunReport], list[str]]] | None = None, 

148) -> tuple[list[WmsRunReport], list[str]]: 

149 """Retrieve summary of jobs submitted for execution. 

150 

151 Parameters 

152 ---------- 

153 wms_service_fqn : `str` 

154 Name of the WMS service class. 

155 run_id : `str`, optional 

156 A run id the report will be restricted to. 

157 user : `str`, optional 

158 A username the report will be restricted to. 

159 hist : `float`, optional 

160 Include runs from the given number of past days. 

161 pass_thru : `str`, optional 

162 A string to pass directly to the WMS service class. 

163 is_global : `bool`, optional 

164 If set, all available job queues will be queried for job information. 

165 Defaults to False which means that only a local job queue will be 

166 queried for information. 

167 

168 Only applicable in the context of a WMS using distributed job queues 

169 (e.g., HTCondor). 

170 return_exit_codes : `bool`, optional 

171 If set, return exit codes related to jobs with a 

172 non-success status. Defaults to False, which means that only 

173 the summary state is returned. 

174 

175 Only applicable in the context of a WMS with associated 

176 handlers to return exit codes from jobs. 

177 postprocessors : `collections.abc.Sequence` [callable], optional 

178 List of functions for "massaging" reports returned by the plugin. Each 

179 function must take one positional argument: 

180 

181 - ``report``: run report (`lsst.ctrl.bps.WmsRunReport`) 

182 

183 If None (default), each run report returned by the plugin (if any) 

184 will be returned as is. 

185 

186 Returns 

187 ------- 

188 reports : `list` [`WmsRunReport`] 

189 Run reports satisfying the search criteria. 

190 messages : `list` [`str`] 

191 Errors that happened during report retrieval and/or processing. 

192 Empty if no issues were encountered. 

193 

194 Raises 

195 ------ 

196 TypeError 

197 Raised if the WMS service class is not a subclass of BaseWmsService. 

198 """ 

199 messages: list[str] = [] 

200 

201 wms_service_class = doImportType(wms_service_fqn) 

202 if not issubclass(wms_service_class, BaseWmsService): 

203 raise TypeError( 

204 f"Invalid WMS service class '{wms_service_fqn}'; must be a subclass of BaseWmsService" 

205 ) 

206 wms_service = wms_service_class({}) 

207 

208 reports, message = wms_service.report( 

209 wms_workflow_id=run_id, 

210 user=user, 

211 hist=hist, 

212 pass_thru=pass_thru, 

213 is_global=is_global, 

214 return_exit_codes=return_exit_codes, 

215 ) 

216 if message: 216 ↛ 217line 216 didn't jump to line 217 because the condition on line 216 was never true

217 messages.append(message) 

218 

219 if postprocessors: 

220 for report in reports: 

221 for postprocessor in postprocessors: 

222 if warnings := postprocessor(report): 

223 for warning in warnings: 

224 messages.append( 

225 f"WARNING: Report may be incomplete. " 

226 f"There was an issue with report postprocessing for '{report.wms_id}': " 

227 f"{warning} (origin: {postprocessor.__name__})" 

228 ) 

229 

230 return reports, messages