Coverage for python/lsst/ctrl/bps/wms_service.py: 99%

141 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-16 09:18 +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"""Base classes for working with a specific WMS.""" 

29 

30__all__ = [ 

31 "BaseWmsService", 

32 "BaseWmsWorkflow", 

33 "WmsJobReport", 

34 "WmsRunReport", 

35 "WmsSpecificInfo", 

36 "WmsStates", 

37] 

38 

39 

40import dataclasses 

41import logging 

42from abc import ABCMeta, abstractmethod 

43from enum import Enum 

44from typing import Any 

45 

46from . import BpsConfig 

47 

48_LOG = logging.getLogger(__name__) 

49 

50 

51class WmsStates(Enum): 

52 """Run and job states.""" 

53 

54 # Offset values so can use as exit codes to bps status 

55 # without colliding with click exit codes (e.g., 2 for 

56 # bad command line) 

57 

58 UNKNOWN = 10 

59 """Can't determine state.""" 

60 

61 MISFIT = 11 

62 """Determined state, but doesn't fit other states.""" 

63 

64 UNREADY = 12 

65 """Still waiting for parents to finish.""" 

66 

67 READY = 13 

68 """All of its parents have finished successfully.""" 

69 

70 PENDING = 14 

71 """Ready to run, visible in batch queue.""" 

72 

73 RUNNING = 15 

74 """Currently running.""" 

75 

76 DELETED = 16 

77 """In the process of being deleted or already deleted.""" 

78 

79 HELD = 17 

80 """In a hold state.""" 

81 

82 SUCCEEDED = 0 

83 """Have completed with success status.""" 

84 

85 FAILED = 19 

86 """Have completed with non-success status.""" 

87 

88 PRUNED = 20 

89 """At least one of the parents failed or can't be run.""" 

90 

91 

92class WmsSpecificInfo: 

93 """Class representing WMS specific information. 

94 

95 Each piece of information is split into two parts: a template and 

96 a context. The template is a string that can contain literal text and/or 

97 *named* replacement fields delimited by braces ``{}``. The context is 

98 a mapping between the names, corresponding to the replacement fields 

99 in the template, and their values. 

100 

101 To produce a human-readable representation of the information, e.g., for 

102 logging purposes, it needs to be rendered first to combine these two parts. 

103 On the other hand, the context alone might be sufficient if the provided 

104 information is being ingested to a database. 

105 """ 

106 

107 def __init__(self) -> None: 

108 self._context: dict[str, Any] = {} 

109 self._templates: list[str] = [] 

110 

111 def __bool__(self) -> bool: 

112 return bool(self._templates) 

113 

114 def __str__(self) -> str: 

115 lines = [] 

116 for template in self._templates: 

117 lines.append(template.format_map(self._context)) 

118 return "\n".join(lines) 

119 

120 @property 

121 def context(self) -> dict[str, Any]: 

122 """The context that will be used to render the information. 

123 

124 Returns 

125 ------- 

126 context : `dict` [`str`, `~typing.Any`] 

127 A copy of the dictionary representing the mapping between 

128 *every* template variable and its value. 

129 

130 Notes 

131 ----- 

132 The property returns a *shallow* copy of the dictionary representing 

133 the context as the intended purpose of the `WmsSpecificInfo` is to 

134 pass a small number of brief messages from WMS to BPS reporting 

135 subsystem. Hence, it is assumed that the dictionary will only contain 

136 immutable objects (e.g. strings, numbers). 

137 """ 

138 return self._context.copy() 

139 

140 @property 

141 def templates(self) -> list[str]: 

142 """The list of templates that will be used to render the information. 

143 

144 Returns 

145 ------- 

146 templates : `list` [`str`] 

147 A copy of the complete list of the message templates in order 

148 in which the messages were added. 

149 """ 

150 return self._templates.copy() 

151 

152 def add_message(self, template: str, context: dict[str, Any] | None = None, **kwargs) -> None: 

153 """Add a message to the WMS information. 

154 

155 If keyword arguments are specified, the passed context is then updated 

156 with those key/value pairs. 

157 

158 Parameters 

159 ---------- 

160 template : `str` 

161 A message template. 

162 context : `dict` [`str`, `~typing.Any`], optional 

163 A mapping between template variables and their values. 

164 **kwargs 

165 Additional keyword arguments. 

166 

167 Raises 

168 ------ 

169 ValueError 

170 Raised if the message can't be rendered due to errors in either 

171 the template, the context, or both. 

172 """ 

173 ctx: dict[str, Any] = {} 

174 if context is not None: 

175 ctx |= context 

176 ctx.update(kwargs) 

177 

178 # Test that given context has all of the values needed for the given 

179 # template. 

180 try: 

181 template.format_map(ctx) 

182 except Exception as exc: 

183 raise ValueError(f"Adding template '{template}' with context '{ctx}' failed") from exc 

184 

185 # Check if the given context does not change values of the already 

186 # existing fields. 

187 common_fields = set(self._context) & set(ctx) 

188 conflicts = [field for field in common_fields if self._context[field] != ctx[field]] 

189 if conflicts: 

190 raise ValueError( 

191 f"Adding template '{template}' with context '{ctx}' failed:" 

192 f"change of value detected for field(s): {', '.join(conflicts)}" 

193 ) 

194 

195 self._context.update(ctx) 

196 self._templates.append(template) 

197 

198 

199@dataclasses.dataclass(slots=True) 

200class WmsJobReport: 

201 """WMS job information to be included in detailed report output.""" 

202 

203 wms_id: str 

204 """Job id assigned by the workflow management system.""" 

205 

206 name: str 

207 """A name assigned automatically by BPS.""" 

208 

209 label: str 

210 """A user-facing label for a job. Multiple jobs can have the same label.""" 

211 

212 state: WmsStates 

213 """Job's current execution state.""" 

214 

215 

216@dataclasses.dataclass(slots=True) 

217class WmsRunReport: 

218 """WMS run information to be included in detailed report output.""" 

219 

220 wms_id: str | None = None 

221 """Id assigned to the run by the WMS. 

222 """ 

223 

224 global_wms_id: str | None = None 

225 """Global run identification number. 

226 

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

228 (e.g., HTCondor). 

229 """ 

230 

231 path: str | None = None 

232 """Path to the submit directory.""" 

233 

234 label: str | None = None 

235 """Run's label.""" 

236 

237 run: str | None = None 

238 """Run's name.""" 

239 

240 project: str | None = None 

241 """Name of the project run belongs to.""" 

242 

243 campaign: str | None = None 

244 """Name of the campaign the run belongs to.""" 

245 

246 payload: str | None = None 

247 """Name of the payload.""" 

248 

249 operator: str | None = None 

250 """Username of the operator who submitted the run.""" 

251 

252 site: str | None = None 

253 """Compute site for payload jobs.""" 

254 

255 run_summary: str | None = None 

256 """Job counts per label.""" 

257 

258 state: WmsStates | None = None 

259 """Run's execution state.""" 

260 

261 jobs: list[WmsJobReport] | None = None 

262 """Information about individual jobs in the run.""" 

263 

264 total_number_jobs: int | None = None 

265 """Total number of jobs in the run.""" 

266 

267 job_state_counts: dict[WmsStates, int] | None = None 

268 """Job counts per state.""" 

269 

270 job_summary: dict[str, dict[WmsStates, int]] | None = None 

271 """Job counts per label and per state.""" 

272 

273 exit_code_summary: dict[str, list[int]] | None = None 

274 """Summary of non-zero exit codes per job label available through the WMS. 

275 

276 Currently behavior for jobs that were canceled, held, etc. are plugin 

277 dependent. 

278 """ 

279 

280 specific_info: WmsSpecificInfo | None = None 

281 """Any additional WMS specific information.""" 

282 

283 

284class BaseWmsService: 

285 """Interface for interactions with a specific WMS. 

286 

287 Parameters 

288 ---------- 

289 config : `lsst.ctrl.bps.BpsConfig` 

290 Configuration needed by the WMS service. 

291 """ 

292 

293 def __init__(self, config): 

294 self.config = config 

295 

296 @property 

297 def defaults(self): 

298 """Service default settings (`lsst.daf.butler.Config`). 

299 

300 Notes 

301 ----- 

302 This property is currently being used in ``BpsConfig.__init__()``. 

303 As long as that's the case it cannot be changed to return 

304 a `BpsConfig` instance. 

305 """ 

306 return None 

307 

308 @property 

309 def defaults_uri(self): 

310 """URI to WMS default settings (`lsst.resources.ResourcePath`).""" 

311 return None 

312 

313 def prepare(self, config, generic_workflow, out_prefix=None): 

314 """Create submission for a generic workflow for a specific WMS. 

315 

316 Parameters 

317 ---------- 

318 config : `lsst.ctrl.bps.BpsConfig` 

319 BPS configuration. 

320 generic_workflow : `lsst.ctrl.bps.GenericWorkflow` 

321 Generic representation of a single workflow. 

322 out_prefix : `str` 

323 Prefix for all WMS output files. 

324 

325 Returns 

326 ------- 

327 wms_workflow : `lsst.ctrl.bps.BaseWmsWorkflow` 

328 Prepared WMS Workflow to submit for execution. 

329 """ 

330 raise NotImplementedError 

331 

332 def submit(self, workflow, **kwargs): 

333 """Submit a single WMS workflow. 

334 

335 Parameters 

336 ---------- 

337 workflow : `lsst.ctrl.bps.BaseWmsWorkflow` 

338 Prepared WMS Workflow to submit for execution. 

339 **kwargs : `~typing.Any` 

340 Additional modifiers to the configuration. 

341 """ 

342 raise NotImplementedError 

343 

344 def restart(self, wms_workflow_id): 

345 """Restart a workflow from the point of failure. 

346 

347 Parameters 

348 ---------- 

349 wms_workflow_id : `str` 

350 Id that can be used by WMS service to identify workflow that 

351 need to be restarted. 

352 

353 Returns 

354 ------- 

355 wms_id : `str` 

356 Id of the restarted workflow. If restart failed, it will be set 

357 to `None`. 

358 run_name : `str` 

359 Name of the restarted workflow. If restart failed, it will be set 

360 to `None`. 

361 message : `str` 

362 A message describing any issues encountered during the restart. 

363 If there were no issue, an empty string is returned. 

364 """ 

365 raise NotImplementedError 

366 

367 def list_submitted_jobs(self, wms_id=None, user=None, require_bps=True, pass_thru=None, is_global=False): 

368 """Query WMS for list of submitted WMS workflows/jobs. 

369 

370 This should be a quick lookup function to create list of jobs for 

371 other functions. 

372 

373 Parameters 

374 ---------- 

375 wms_id : `int` or `str`, optional 

376 Id or path that can be used by WMS service to look up job. 

377 user : `str`, optional 

378 User whose submitted jobs should be listed. 

379 require_bps : `bool`, optional 

380 Whether to require jobs returned in list to be bps-submitted jobs. 

381 pass_thru : `str`, optional 

382 Information to pass through to WMS. 

383 is_global : `bool`, optional 

384 If set, all available job queues will be queried for job 

385 information. Defaults to False which means that only a local job 

386 queue will be queried for information. 

387 

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

389 queues (e.g., HTCondor). A WMS with a centralized job queue 

390 (e.g. PanDA) can safely ignore it. 

391 

392 Returns 

393 ------- 

394 job_ids : `list` [`~typing.Any`] 

395 Only job ids to be used by cancel and other functions. Typically 

396 this means top-level jobs (i.e., not children jobs). 

397 """ 

398 raise NotImplementedError 

399 

400 def report( 

401 self, 

402 wms_workflow_id=None, 

403 user=None, 

404 hist=0, 

405 pass_thru=None, 

406 is_global=False, 

407 return_exit_codes=False, 

408 ): 

409 """Query WMS for status of submitted WMS workflows. 

410 

411 Parameters 

412 ---------- 

413 wms_workflow_id : `int` or `str`, optional 

414 Id that can be used by WMS service to look up status. 

415 user : `str`, optional 

416 Limit report to submissions by this particular user. 

417 hist : `float`, optional 

418 Number of days to expand report to include finished WMS workflows. 

419 pass_thru : `str`, optional 

420 Additional arguments to pass through to the specific WMS service. 

421 is_global : `bool`, optional 

422 If set, all available job queues will be queried for job 

423 information. Defaults to False which means that only a local job 

424 queue will be queried for information. 

425 

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

427 queues (e.g., HTCondor). A WMS with a centralized job queue 

428 (e.g. PanDA) can safely ignore it. 

429 return_exit_codes : `bool`, optional 

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

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

432 the summary state is returned. 

433 

434 Only applicable in the context of a WMS with associated 

435 handlers to return exit codes from jobs. 

436 

437 Returns 

438 ------- 

439 run_reports : `list` [`lsst.ctrl.bps.WmsRunReport`] 

440 Status information for submitted WMS workflows. 

441 message : `str` 

442 Message to user on how to find more status information specific to 

443 this particular WMS. 

444 """ 

445 raise NotImplementedError 

446 

447 def get_status( 

448 self, 

449 wms_workflow_id: str, 

450 hist: float = 1, 

451 is_global: bool = False, 

452 ) -> tuple[WmsStates, str]: 

453 """Query WMS for quick status of single submitted WMS workflow. 

454 

455 Parameters 

456 ---------- 

457 wms_workflow_id : `int` or `str`, optional 

458 ID that can be used by WMS service to look up status. 

459 hist : `float`, optional 

460 Number of days to expand query to include finished WMS workflows. 

461 Defaults to 1. 

462 is_global : `bool`, optional 

463 If set, all available job queues will be queried for run 

464 information. Defaults to False which means that only a local run 

465 queue will be queried for information. 

466 

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

468 queues (e.g., HTCondor). A WMS with a centralized job queue 

469 (e.g. PanDA) can safely ignore it. 

470 

471 Returns 

472 ------- 

473 status : `lsst.ctrl.bps.WmsStates` 

474 Status of single run from given information. 

475 message : `str` 

476 Extra message for status command to print. This could be pointers 

477 to documentation or to WMS specific commands. 

478 """ 

479 raise NotImplementedError 

480 

481 def cancel(self, wms_id, pass_thru=None): 

482 """Cancel submitted workflows/jobs. 

483 

484 Parameters 

485 ---------- 

486 wms_id : `str` 

487 ID or path of job that should be canceled. 

488 pass_thru : `str`, optional 

489 Information to pass through to WMS. 

490 

491 Returns 

492 ------- 

493 deleted : `bool` 

494 Whether successful deletion or not. Currently, if any doubt or any 

495 individual jobs not deleted, return False. 

496 message : `str` 

497 Any message from WMS (e.g., error details). 

498 """ 

499 raise NotImplementedError 

500 

501 def run_submission_checks(self): 

502 """Check to run at start if running WMS specific submission steps. 

503 

504 Any exception other than NotImplementedError will halt submission. 

505 Submit directory may not yet exist when this is called. 

506 """ 

507 raise NotImplementedError 

508 

509 def ping(self, pass_thru): 

510 """Check whether WMS services are up, reachable, and can authenticate 

511 if authentication is required. 

512 

513 The services to be checked are those needed for submit, report, cancel, 

514 restart, but ping cannot guarantee whether jobs would actually run 

515 successfully. 

516 

517 Parameters 

518 ---------- 

519 pass_thru : `str`, optional 

520 Information to pass through to WMS. 

521 

522 Returns 

523 ------- 

524 status : `int` 

525 0 for success, non-zero for failure. 

526 message : `str` 

527 Any message from WMS (e.g., error details). 

528 """ 

529 raise NotImplementedError 

530 

531 

532class BaseWmsWorkflow(metaclass=ABCMeta): 

533 """Interface for single workflow specific to a WMS. 

534 

535 Parameters 

536 ---------- 

537 name : `str` 

538 Unique name of workflow. 

539 config : `lsst.ctrl.bps.BpsConfig` 

540 Generic workflow config. 

541 """ 

542 

543 def __init__(self, name, config): 

544 self.name = name 

545 self.config = config 

546 self.service_class = None 

547 self.run_id = None 

548 self.submit_path = None 

549 

550 @classmethod 

551 def from_generic_workflow(cls, config, generic_workflow, out_prefix, service_class): 

552 """Create a WMS-specific workflow from a GenericWorkflow. 

553 

554 Parameters 

555 ---------- 

556 config : `lsst.ctrl.bps.BpsConfig` 

557 Configuration values needed for generating a WMS specific workflow. 

558 generic_workflow : `lsst.ctrl.bps.GenericWorkflow` 

559 Generic workflow from which to create the WMS-specific one. 

560 out_prefix : `str` 

561 Root directory to be used for WMS workflow inputs and outputs 

562 as well as internal WMS files. 

563 service_class : `str` 

564 Full module name of WMS service class that created this workflow. 

565 

566 Returns 

567 ------- 

568 wms_workflow : `lsst.ctrl.bps.BaseWmsWorkflow` 

569 A WMS specific workflow. 

570 """ 

571 raise NotImplementedError 

572 

573 @abstractmethod 

574 def write(self, out_prefix): 

575 """Write WMS files for this particular workflow. 

576 

577 Parameters 

578 ---------- 

579 out_prefix : `str` 

580 Root directory to be used for WMS workflow inputs and outputs 

581 as well as internal WMS files. 

582 """ 

583 raise NotImplementedError 

584 

585 def add_to_parent_workflow(self, config: BpsConfig) -> None: 

586 """Add self to parent workflow. 

587 

588 Parameters 

589 ---------- 

590 config : `lsst.ctrl.bps.BpsConfig` 

591 BPS configuration. 

592 """ 

593 raise NotImplementedError