Coverage for python/lsst/pipe/base/_status.py: 92%
141 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-19 09:26 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-19 09:26 +0000
1# This file is part of pipe_base.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (http://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 <http://www.gnu.org/licenses/>.
28from __future__ import annotations
30__all__ = (
31 "AlgorithmError",
32 "AnnotatedPartialOutputsError",
33 "ExceptionInfo",
34 "InvalidQuantumError",
35 "NoWorkFound",
36 "QuantumAttemptStatus",
37 "QuantumSuccessCaveats",
38 "RepeatableQuantumError",
39 "UnprocessableDataError",
40 "UpstreamFailureNoWorkFound",
41)
43import abc
44import enum
45import logging
46import re
47import sys
48from collections.abc import MutableMapping
49from typing import TYPE_CHECKING, Any, ClassVar, Protocol
51import pydantic
53from lsst.utils import introspection
54from lsst.utils.logging import LsstLogAdapter, getLogger
56from ._task_metadata import GetSetDictMetadata, NestedMetadataDict
58if TYPE_CHECKING:
59 from ._task_metadata import TaskMetadata
62_LOG = getLogger(__name__)
65class QuantumSuccessCaveats(enum.Flag):
66 """Flags that add caveats to a "successful" quantum.
68 Quanta can be considered successful even if they do not produce some of
69 their expected outputs (and even if they do not produce all of their
70 expected outputs), as long as the condition is sufficiently well understood
71 that downstream processing should succeed.
72 """
74 NO_CAVEATS = 0
75 """All outputs were produced and no exceptions were raised."""
77 ANY_OUTPUTS_MISSING = enum.auto()
78 """At least one predicted output was not produced."""
80 ALL_OUTPUTS_MISSING = enum.auto()
81 """No predicted outputs (except logs and metadata) were produced.
83 `ANY_OUTPUTS_MISSING` is also set whenever this flag is set.
84 """
86 NO_WORK = enum.auto()
87 """A subclass of `NoWorkFound` was raised.
89 This does not necessarily imply that `ANY_OUTPUTS_MISSING` is not set,
90 since a `PipelineTask.runQuantum` implementation could raise it after
91 directly writing all of its predicted outputs.
92 """
94 ADJUST_QUANTUM_RAISED = enum.auto()
95 """`NoWorkFound` was raised by `PipelineTaskConnnections.adjustQuantum`.
97 This indicates that if a new `QuantumGraph` had been generated immediately
98 before running this quantum, that quantum would not have even been
99 included, because required inputs that were expected to exist by the time
100 it was run (in the original `QuantumGraph`) were not actually produced.
102 `NO_WORK` and `ALL_OUTPUTS_MISSING` are also set whenever this flag is set.
103 """
105 UPSTREAM_FAILURE_NO_WORK = enum.auto()
106 """`UpstreamFailureNoWorkFound` was raised by `PipelineTask.runQuantum`.
108 This exception is raised by downstream tasks when an upstream task's
109 outputs were incomplete in a way that blocks it from running, often
110 because the upstream task raised `AnnotatedPartialOutputsError`.
112 `NO_WORK` is also set whenever this flag is set.
113 """
115 UNPROCESSABLE_DATA = enum.auto()
116 """`UnprocessableDataError` was raised by `PipelineTask.runQuantum`.
118 `NO_WORK` is also set whenever this flag is set.
119 """
121 PARTIAL_OUTPUTS_ERROR = enum.auto()
122 """`AnnotatedPartialOutputsError` was raised by `PipelineTask.runQuantum`
123 and the execution system was instructed to consider this a qualified
124 success.
125 """
127 @classmethod
128 def from_adjust_quantum_no_work(cls) -> QuantumSuccessCaveats:
129 """Return the set of flags appropriate for a quantum for which
130 `PipelineTaskConnections.adjustdQuantum` raised `NoWorkFound`.
131 """
132 return cls.NO_WORK | cls.ADJUST_QUANTUM_RAISED | cls.ANY_OUTPUTS_MISSING | cls.ALL_OUTPUTS_MISSING
134 @classmethod
135 def expanded_dict(cls, concise: str) -> dict | None:
136 """Return a dictionary representation of the concise flags.
138 Parameters
139 ----------
140 concise : `str`
141 The concise string representation of the flags to expand to a
142 dictionary.
144 Returns
145 -------
146 d : `dict` | `None`
147 A dictionary expansion of the concise flag string with `token`,
148 `code`, and `count` keys; or `None` if the string does not
149 represent a concise flag.
150 """
151 r = re.compile(r"(?P<token>\*|\+)?(?P<code>[A-Z]{1})\((?P<count>[0-9]+)\)")
152 m = re.match(r, concise)
153 return m.groupdict() if m is not None else None
155 def concise(self) -> str:
156 """Return a concise string representation of the flags.
158 Returns
159 -------
160 s : `str`
161 Two-character string representation, with the first character
162 indicating whether any predicted outputs were missing and the
163 second representing any exceptions raised. This representation is
164 not always complete; some rare combinations of flags are displayed
165 as if only one of the flags was set.
167 Notes
168 -----
169 The `legend` method returns a description of the returned codes.
170 """
171 char1 = ""
172 if self & QuantumSuccessCaveats.ALL_OUTPUTS_MISSING:
173 char1 = "*"
174 elif self & QuantumSuccessCaveats.ANY_OUTPUTS_MISSING:
175 char1 = "+"
176 char2 = ""
177 if self & QuantumSuccessCaveats.ADJUST_QUANTUM_RAISED:
178 char2 = "A"
179 elif self & QuantumSuccessCaveats.UNPROCESSABLE_DATA:
180 char2 = "D"
181 elif self & QuantumSuccessCaveats.UPSTREAM_FAILURE_NO_WORK:
182 char2 = "U"
183 elif self & QuantumSuccessCaveats.PARTIAL_OUTPUTS_ERROR:
184 char2 = "P"
185 elif self & QuantumSuccessCaveats.NO_WORK: 185 ↛ 187line 185 didn't jump to line 187 because the condition on line 185 was always true
186 char2 = "N"
187 return char1 + char2
189 @staticmethod
190 def legend() -> dict[str, str]:
191 """Return a `dict` with human-readable descriptions of the characters
192 used in `concise`.
194 Returns
195 -------
196 legend : `dict` [ `str`, `str` ]
197 Mapping from character code to description.
198 """
199 return {
200 "+": "at least one predicted output was missing, but not all were",
201 "*": "all predicted outputs were missing (besides logs and metadata)",
202 "A": "adjustQuantum raised NoWorkFound; a regenerated QG would not include this quantum",
203 "D": "algorithm considers data too bad to be processable",
204 "U": "one or more input dataset was incomplete due to an upstream failure",
205 "P": "task failed but wrote partial outputs; considered a partial success",
206 "N": "runQuantum raised NoWorkFound",
207 }
210class ExceptionInfo(pydantic.BaseModel):
211 """Information about an exception that was raised."""
213 type_name: str
214 """Fully-qualified Python type name for the exception raised."""
216 message: str
217 """String message included in the exception."""
219 metadata: dict[str, float | int | str | bool | None]
220 """Additional metadata included in the exception."""
222 @classmethod
223 def _from_metadata(cls, md: TaskMetadata) -> ExceptionInfo:
224 """Construct from task metadata.
226 Parameters
227 ----------
228 md : `TaskMetadata`
229 Metadata about the error, as written by
230 `AnnotatedPartialOutputsError`.
232 Returns
233 -------
234 info : `ExceptionInfo`
235 Information about the exception.
236 """
237 result = cls(type_name=md["type"], message=md["message"], metadata={})
238 if "metadata" in md: 238 ↛ 250line 238 didn't jump to line 250 because the condition on line 238 was always true
239 raw_err_metadata = md["metadata"].to_dict()
240 for k, v in raw_err_metadata.items():
241 # Guard against error metadata we wouldn't be able to serialize
242 # later via Pydantic; don't want one weird value bringing down
243 # our ability to report on an entire run.
244 if isinstance(v, float | int | str | bool): 244 ↛ 247line 244 didn't jump to line 247 because the condition on line 244 was always true
245 result.metadata[k] = v
246 else:
247 _LOG.debug(
248 "Not propagating nested or JSON-incompatible exception metadata key %s=%r.", k, v
249 )
250 return result
252 # Work around the fact that Sphinx chokes on Pydantic docstring formatting,
253 # when we inherit those docstrings in our public classes.
254 if "sphinx" in sys.modules and not TYPE_CHECKING:
256 def copy(self, *args: Any, **kwargs: Any) -> Any:
257 """See `pydantic.BaseModel.copy`."""
258 return super().copy(*args, **kwargs)
260 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
261 """See `pydantic.BaseModel.model_dump`."""
262 return super().model_dump(*args, **kwargs)
264 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any:
265 """See `pydantic.BaseModel.model_dump_json`."""
266 return super().model_dump(*args, **kwargs)
268 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
269 """See `pydantic.BaseModel.model_copy`."""
270 return super().model_copy(*args, **kwargs)
272 @classmethod
273 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override]
274 """See `pydantic.BaseModel.model_construct`."""
275 return super().model_construct(*args, **kwargs)
277 @classmethod
278 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any:
279 """See `pydantic.BaseModel.model_json_schema`."""
280 return super().model_json_schema(*args, **kwargs)
282 @classmethod
283 def model_validate(cls, *args: Any, **kwargs: Any) -> Any:
284 """See `pydantic.BaseModel.model_validate`."""
285 return super().model_validate(*args, **kwargs)
287 @classmethod
288 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any:
289 """See `pydantic.BaseModel.model_validate_json`."""
290 return super().model_validate_json(*args, **kwargs)
292 @classmethod
293 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any:
294 """See `pydantic.BaseModel.model_validate_strings`."""
295 return super().model_validate_strings(*args, **kwargs)
298class QuantumAttemptStatus(enum.Enum):
299 """Enum summarizing an attempt to run a quantum."""
301 ABORTED = -4
302 """The quantum failed with a hard error that prevented both logs and
303 metadata from being written.
305 This state is only set if information from higher-level tooling (e.g. BPS)
306 is available to distinguish it from ``UNKNOWN``.
307 """
309 UNKNOWN = -3
310 """The status of this attempt is unknown.
312 This means no logs or metadata were written, and it at least could not be
313 determined whether the quantum was blocked by an upstream failure (if it
314 was definitely blocked, `BLOCKED` is set instead).
315 """
317 ABORTED_SUCCESS = -2
318 """Task metadata was written for this attempt but logs were not.
320 This is a rare condition that requires a hard failure (i.e. the kind that
321 can prevent a ``finally`` block from running or I/O from being durable) at
322 a very precise time.
323 """
325 FAILED = -1
326 """Execution of the quantum failed gracefully.
328 This is always set if the task metadata dataset was not written but logs
329 were, as is the case when a Python exception is caught and handled by the
330 execution system.
332 This status guarantees that the task log dataset was produced but the
333 metadata dataset was not.
334 """
336 BLOCKED = 0
337 """This quantum was not executed because an upstream quantum failed.
339 Upstream quanta with status `UNKNOWN`, `FAILED`, or `ABORTED` are
340 considered blockers; `ABORTED_SUCCESS` is not.
341 """
343 SUCCESSFUL = 1
344 """This quantum was successfully executed.
346 Quanta may be considered successful even if they do not write any outputs
347 or shortcut early by raising `NoWorkFound` or one of its variants. They
348 may even be considered successful if they raise
349 `AnnotatedPartialOutputsError` if the executor is configured to treat that
350 exception as a non-failure. See `QuantumSuccessCaveats` for details on how
351 these "successes with caveats" are reported.
352 """
354 @property
355 def has_metadata(self) -> bool:
356 """Whether the task metadata dataset was produced."""
357 return self is self.SUCCESSFUL or self is self.ABORTED_SUCCESS
359 @property
360 def has_log(self) -> bool:
361 """Whether the log dataset was produced."""
362 return self is self.SUCCESSFUL or self is self.FAILED
364 @property
365 def title(self) -> str:
366 """A version of this status' name suitable for use as a title in a plot
367 or table.
368 """
369 return self.name.capitalize().replace("_", " ")
371 @property
372 def is_rare(self) -> bool:
373 """Whether this status is rare enough that it should only be listed
374 when it actually occurs.
375 """
376 return self in (self.ABORTED, self.ABORTED_SUCCESS, self.UNKNOWN)
379class GetSetDictMetadataHolder(Protocol):
380 """Protocol for objects that have a ``metadata`` attribute that satisfies
381 `GetSetDictMetadata`.
382 """
384 @property
385 def metadata(self) -> GetSetDictMetadata | MutableMapping[str, Any] | None:
386 pass
389class NoWorkFound(BaseException):
390 """An exception raised when a Quantum should not exist because there is no
391 work for it to do.
393 This usually occurs because a non-optional input dataset is not present, or
394 a spatiotemporal overlap that was conservatively predicted does not
395 actually exist.
397 This inherits from BaseException because it is used to signal a case that
398 we don't consider a real error, even though we often want to use try/except
399 logic to trap it.
400 """
402 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK
405class UpstreamFailureNoWorkFound(NoWorkFound):
406 """A specialization of `NoWorkFound` that indicates that an upstream task
407 had a problem that was ignored (e.g. to prevent a single-detector failure
408 from bringing down an entire visit).
409 """
411 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK | QuantumSuccessCaveats.UPSTREAM_FAILURE_NO_WORK
414class RepeatableQuantumError(RuntimeError):
415 """Exception that may be raised by PipelineTasks (and code they delegate
416 to) in order to indicate that a repeatable problem that will not be
417 addressed by retries.
419 This usually indicates that the algorithm and the data it has been given
420 are somehow incompatible, and the task should run fine on most other data.
422 This exception may be used as a base class for more specific questions, or
423 used directly while chaining another exception, e.g.::
425 try:
426 run_code()
427 except SomeOtherError as err:
428 raise RepeatableQuantumError() from err
430 This may be used for missing input data when the desired behavior is to
431 cause all downstream tasks being run be blocked, forcing the user to
432 address the problem. When the desired behavior is to skip all of this
433 quantum and attempt downstream tasks (or skip them) without its its
434 outputs, raise `NoWorkFound` or return without raising instead.
435 """
437 EXIT_CODE = 20
440class AlgorithmError(RepeatableQuantumError, abc.ABC):
441 """Exception that may be raised by PipelineTasks (and code they delegate
442 to) in order to indicate a repeatable algorithmic failure that will not be
443 addressed by retries.
445 Subclass this exception to define the metadata associated with the error
446 (for example: number of data points in a fit vs. degrees of freedom).
447 """
449 def __new__(cls, *args: Any, **kwargs: Any) -> AlgorithmError:
450 # Have to override __new__ because builtin subclasses aren't checked
451 # for abstract methods; see https://github.com/python/cpython/issues/50246
452 if cls.__abstractmethods__:
453 raise TypeError(
454 f"Can't instantiate abstract class {cls.__name__} with "
455 f"abstract methods: {','.join(sorted(cls.__abstractmethods__))}"
456 )
457 return super().__new__(cls, *args, **kwargs)
459 @property
460 @abc.abstractmethod
461 def metadata(self) -> NestedMetadataDict | None:
462 """Metadata from the raising `~lsst.pipe.base.Task` with more
463 information about the failure. The contents of the dict are
464 `~lsst.pipe.base.Task`-dependent, and must have `str` keys and `str`,
465 `int`, `float`, `bool`, or nested-dictionary (with the same key and
466 value types) values.
467 """
468 raise NotImplementedError
471class UnprocessableDataError(NoWorkFound):
472 """A specialization of `NoWorkFound` that will be [subclassed and] raised
473 by Tasks to indicate a failure to process their inputs for some reason that
474 is non-recoverable.
476 Notes
477 -----
478 An example is a known bright star that causes PSF measurement to fail, and
479 that makes that detector entirely non-recoverable. Another example is an
480 image with an oddly shaped PSF (e.g. due to a failure to achieve focus)
481 that warrants the image being flagged as "poor quality" which should not
482 have further processing attempted.
484 The `NoWorkFound` inheritance ensures the job will not be considered a
485 failure (i.e. such that no human time will inadvertently be spent chasing
486 it down).
488 Do not raise this unless we are convinced that the data cannot (or should
489 not) be processed, even by a better algorithm. Most instances where this
490 error would be raised likely require an RFC to explicitly define the
491 situation.
492 """
494 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK | QuantumSuccessCaveats.UNPROCESSABLE_DATA
497class AnnotatedPartialOutputsError(RepeatableQuantumError):
498 """Exception that runQuantum raises when the (partial) outputs it has
499 written contain information about their own incompleteness or degraded
500 quality.
502 Clients should construct this exception by calling `annotate` instead of
503 calling the constructor directly. However, `annotate` does not chain the
504 exception; this must still be done by the client.
506 This exception should always chain the original error. When the
507 executor catches this exception, it will report the original exception. In
508 contrast, other exceptions raised from ``runQuantum`` are considered to
509 invalidate any outputs that are already written.
510 """
512 FLAGS: ClassVar = QuantumSuccessCaveats.PARTIAL_OUTPUTS_ERROR
514 @classmethod
515 def annotate(
516 cls, error: Exception, *args: GetSetDictMetadataHolder | None, log: logging.Logger | LsstLogAdapter
517 ) -> AnnotatedPartialOutputsError:
518 """Set metadata on outputs to explain the nature of the failure.
520 Parameters
521 ----------
522 error : `Exception`
523 Exception that caused the task to fail.
524 *args : `GetSetDictMetadataHolder`
525 Objects (e.g. Task, Exposure, SimpleCatalog) to annotate with
526 failure information. They must have a `metadata` property that
527 is a `~collections.abc.MutableMapping`.
528 log : `logging.Logger`
529 Log to send error message to.
531 Returns
532 -------
533 error : `AnnotatedPartialOutputsError`
534 Exception that the failing task can ``raise from`` with the
535 passed-in exception.
537 Notes
538 -----
539 This should be called from within an except block that has caught an
540 exception. Here is an example of handling a failure in
541 ``PipelineTask.runQuantum`` that annotates and writes partial outputs:
543 .. code-block:: py
544 :name: annotate-error-example
546 def runQuantum(self, butlerQC, inputRefs, outputRefs):
547 inputs = butlerQC.get(inputRefs)
548 exposures = inputs.pop("exposures")
549 assert not inputs, "runQuantum got more inputs than expected"
551 result = pipeBase.Struct(catalog=None)
552 try:
553 self.run(exposure)
554 except pipeBase.AlgorithmError as e:
555 error = pipeBase.AnnotatedPartialOutputsError.annotate(
556 e, self, result.catalog, log=self.log
557 )
558 raise error from e
559 finally:
560 butlerQC.put(result, outputRefs)
561 """
562 failure_info = {
563 "message": str(error),
564 "type": introspection.get_full_type_name(error),
565 }
566 if other := getattr(error, "metadata", None):
567 failure_info["metadata"] = other
569 # NOTE: Can't fully test this in pipe_base because afw is not a
570 # dependency; test_calibrateImage.py in pipe_tasks gives more coverage.
571 for item in args:
572 # Some outputs may not exist, so we cannot set metadata on them.
573 if item is None: 573 ↛ 574line 573 didn't jump to line 574 because the condition on line 573 was never true
574 continue
575 if hasattr(item.metadata, "set_dict"): 575 ↛ 577line 575 didn't jump to line 577 because the condition on line 575 was always true
576 item.metadata.set_dict("failure", failure_info) # type: ignore
577 elif isinstance(item.metadata, MutableMapping):
578 item.metadata["failure"] = failure_info
579 else:
580 raise TypeError(f"Invalid metadata type {type(item.metadata).__name__!r}")
582 log.debug(
583 "Task failed with only partial outputs; see exception message for details.",
584 exc_info=error,
585 )
587 return cls("Task failed and wrote partial outputs: see chained exception for details.")
590class InvalidQuantumError(Exception):
591 """Exception that may be raised by PipelineTasks (and code they delegate
592 to) in order to indicate logic bug or configuration problem.
594 This usually indicates that the configured algorithm itself is invalid and
595 will not run on a significant fraction of quanta (often all of them).
597 This exception may be used as a base class for more specific questions, or
598 used directly while chaining another exception, e.g.::
600 try:
601 run_code()
602 except SomeOtherError as err:
603 raise RepeatableQuantumError() from err
605 Raising this exception in `PipelineTask.runQuantum` or something it calls
606 is a last resort - whenever possible, such problems should cause exceptions
607 in ``__init__`` or in QuantumGraph generation. It should never be used
608 for missing data.
609 """
611 EXIT_CODE = 21