Coverage for python/lsst/ctrl/bps/wms_service.py: 99%
141 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-19 09:28 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-19 09:28 +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/>.
28"""Base classes for working with a specific WMS."""
30__all__ = [
31 "BaseWmsService",
32 "BaseWmsWorkflow",
33 "WmsJobReport",
34 "WmsRunReport",
35 "WmsSpecificInfo",
36 "WmsStates",
37]
40import dataclasses
41import logging
42from abc import ABCMeta, abstractmethod
43from enum import Enum
44from typing import Any
46from . import BpsConfig
48_LOG = logging.getLogger(__name__)
51class WmsStates(Enum):
52 """Run and job states."""
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)
58 UNKNOWN = 10
59 """Can't determine state."""
61 MISFIT = 11
62 """Determined state, but doesn't fit other states."""
64 UNREADY = 12
65 """Still waiting for parents to finish."""
67 READY = 13
68 """All of its parents have finished successfully."""
70 PENDING = 14
71 """Ready to run, visible in batch queue."""
73 RUNNING = 15
74 """Currently running."""
76 DELETED = 16
77 """In the process of being deleted or already deleted."""
79 HELD = 17
80 """In a hold state."""
82 SUCCEEDED = 0
83 """Have completed with success status."""
85 FAILED = 19
86 """Have completed with non-success status."""
88 PRUNED = 20
89 """At least one of the parents failed or can't be run."""
92class WmsSpecificInfo:
93 """Class representing WMS specific information.
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.
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 """
107 def __init__(self) -> None:
108 self._context: dict[str, Any] = {}
109 self._templates: list[str] = []
111 def __bool__(self) -> bool:
112 return bool(self._templates)
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)
120 @property
121 def context(self) -> dict[str, Any]:
122 """The context that will be used to render the information.
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.
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()
140 @property
141 def templates(self) -> list[str]:
142 """The list of templates that will be used to render the information.
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()
152 def add_message(self, template: str, context: dict[str, Any] | None = None, **kwargs) -> None:
153 """Add a message to the WMS information.
155 If keyword arguments are specified, the passed context is then updated
156 with those key/value pairs.
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.
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)
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
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 )
195 self._context.update(ctx)
196 self._templates.append(template)
199@dataclasses.dataclass(slots=True)
200class WmsJobReport:
201 """WMS job information to be included in detailed report output."""
203 wms_id: str
204 """Job id assigned by the workflow management system."""
206 name: str
207 """A name assigned automatically by BPS."""
209 label: str
210 """A user-facing label for a job. Multiple jobs can have the same label."""
212 state: WmsStates
213 """Job's current execution state."""
216@dataclasses.dataclass(slots=True)
217class WmsRunReport:
218 """WMS run information to be included in detailed report output."""
220 wms_id: str | None = None
221 """Id assigned to the run by the WMS.
222 """
224 global_wms_id: str | None = None
225 """Global run identification number.
227 Only applicable in the context of a WMS using distributed job queues
228 (e.g., HTCondor).
229 """
231 path: str | None = None
232 """Path to the submit directory."""
234 label: str | None = None
235 """Run's label."""
237 run: str | None = None
238 """Run's name."""
240 project: str | None = None
241 """Name of the project run belongs to."""
243 campaign: str | None = None
244 """Name of the campaign the run belongs to."""
246 payload: str | None = None
247 """Name of the payload."""
249 operator: str | None = None
250 """Username of the operator who submitted the run."""
252 site: str | None = None
253 """Compute site for payload jobs."""
255 run_summary: str | None = None
256 """Job counts per label."""
258 state: WmsStates | None = None
259 """Run's execution state."""
261 jobs: list[WmsJobReport] | None = None
262 """Information about individual jobs in the run."""
264 total_number_jobs: int | None = None
265 """Total number of jobs in the run."""
267 job_state_counts: dict[WmsStates, int] | None = None
268 """Job counts per state."""
270 job_summary: dict[str, dict[WmsStates, int]] | None = None
271 """Job counts per label and per state."""
273 exit_code_summary: dict[str, list[int]] | None = None
274 """Summary of non-zero exit codes per job label available through the WMS.
276 Currently behavior for jobs that were canceled, held, etc. are plugin
277 dependent.
278 """
280 specific_info: WmsSpecificInfo | None = None
281 """Any additional WMS specific information."""
284class BaseWmsService:
285 """Interface for interactions with a specific WMS.
287 Parameters
288 ----------
289 config : `lsst.ctrl.bps.BpsConfig`
290 Configuration needed by the WMS service.
291 """
293 def __init__(self, config):
294 self.config = config
296 @property
297 def defaults(self):
298 """Service default settings (`lsst.daf.butler.Config`).
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
308 @property
309 def defaults_uri(self):
310 """URI to WMS default settings (`lsst.resources.ResourcePath`)."""
311 return None
313 def prepare(self, config, generic_workflow, out_prefix=None):
314 """Create submission for a generic workflow for a specific WMS.
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.
325 Returns
326 -------
327 wms_workflow : `lsst.ctrl.bps.BaseWmsWorkflow`
328 Prepared WMS Workflow to submit for execution.
329 """
330 raise NotImplementedError
332 def submit(self, workflow, **kwargs):
333 """Submit a single WMS workflow.
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
344 def restart(self, wms_workflow_id):
345 """Restart a workflow from the point of failure.
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.
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
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.
370 This should be a quick lookup function to create list of jobs for
371 other functions.
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.
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.
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
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.
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.
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.
434 Only applicable in the context of a WMS with associated
435 handlers to return exit codes from jobs.
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
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.
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.
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.
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
481 def cancel(self, wms_id, pass_thru=None):
482 """Cancel submitted workflows/jobs.
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.
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
501 def run_submission_checks(self):
502 """Check to run at start if running WMS specific submission steps.
504 Any exception other than NotImplementedError will halt submission.
505 Submit directory may not yet exist when this is called.
506 """
507 raise NotImplementedError
509 def ping(self, pass_thru):
510 """Check whether WMS services are up, reachable, and can authenticate
511 if authentication is required.
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.
517 Parameters
518 ----------
519 pass_thru : `str`, optional
520 Information to pass through to WMS.
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
532class BaseWmsWorkflow(metaclass=ABCMeta):
533 """Interface for single workflow specific to a WMS.
535 Parameters
536 ----------
537 name : `str`
538 Unique name of workflow.
539 config : `lsst.ctrl.bps.BpsConfig`
540 Generic workflow config.
541 """
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
550 @classmethod
551 def from_generic_workflow(cls, config, generic_workflow, out_prefix, service_class):
552 """Create a WMS-specific workflow from a GenericWorkflow.
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.
566 Returns
567 -------
568 wms_workflow : `lsst.ctrl.bps.BaseWmsWorkflow`
569 A WMS specific workflow.
570 """
571 raise NotImplementedError
573 @abstractmethod
574 def write(self, out_prefix):
575 """Write WMS files for this particular workflow.
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
585 def add_to_parent_workflow(self, config: BpsConfig) -> None:
586 """Add self to parent workflow.
588 Parameters
589 ----------
590 config : `lsst.ctrl.bps.BpsConfig`
591 BPS configuration.
592 """
593 raise NotImplementedError