Coverage for python/lsst/pipe/base/_status.py: 92%
135 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-14 07:20 +0000
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-14 07:20 +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 sys
47from collections.abc import MutableMapping
48from typing import TYPE_CHECKING, Any, ClassVar, Protocol
50import pydantic
52from lsst.utils import introspection
53from lsst.utils.logging import LsstLogAdapter, getLogger
55from ._task_metadata import GetSetDictMetadata, NestedMetadataDict
57if TYPE_CHECKING:
58 from ._task_metadata import TaskMetadata
61_LOG = getLogger(__name__)
64class QuantumSuccessCaveats(enum.Flag):
65 """Flags that add caveats to a "successful" quantum.
67 Quanta can be considered successful even if they do not produce some of
68 their expected outputs (and even if they do not produce all of their
69 expected outputs), as long as the condition is sufficiently well understood
70 that downstream processing should succeed.
71 """
73 NO_CAVEATS = 0
74 """All outputs were produced and no exceptions were raised."""
76 ANY_OUTPUTS_MISSING = enum.auto()
77 """At least one predicted output was not produced."""
79 ALL_OUTPUTS_MISSING = enum.auto()
80 """No predicted outputs (except logs and metadata) were produced.
82 `ANY_OUTPUTS_MISSING` is also set whenever this flag is set.
83 """
85 NO_WORK = enum.auto()
86 """A subclass of `NoWorkFound` was raised.
88 This does not necessarily imply that `ANY_OUTPUTS_MISSING` is not set,
89 since a `PipelineTask.runQuantum` implementation could raise it after
90 directly writing all of its predicted outputs.
91 """
93 ADJUST_QUANTUM_RAISED = enum.auto()
94 """`NoWorkFound` was raised by `PipelineTaskConnnections.adjustQuantum`.
96 This indicates that if a new `QuantumGraph` had been generated immediately
97 before running this quantum, that quantum would not have even been
98 included, because required inputs that were expected to exist by the time
99 it was run (in the original `QuantumGraph`) were not actually produced.
101 `NO_WORK` and `ALL_OUTPUTS_MISSING` are also set whenever this flag is set.
102 """
104 UPSTREAM_FAILURE_NO_WORK = enum.auto()
105 """`UpstreamFailureNoWorkFound` was raised by `PipelineTask.runQuantum`.
107 This exception is raised by downstream tasks when an upstream task's
108 outputs were incomplete in a way that blocks it from running, often
109 because the upstream task raised `AnnotatedPartialOutputsError`.
111 `NO_WORK` is also set whenever this flag is set.
112 """
114 UNPROCESSABLE_DATA = enum.auto()
115 """`UnprocessableDataError` was raised by `PipelineTask.runQuantum`.
117 `NO_WORK` is also set whenever this flag is set.
118 """
120 PARTIAL_OUTPUTS_ERROR = enum.auto()
121 """`AnnotatedPartialOutputsError` was raised by `PipelineTask.runQuantum`
122 and the execution system was instructed to consider this a qualified
123 success.
124 """
126 @classmethod
127 def from_adjust_quantum_no_work(cls) -> QuantumSuccessCaveats:
128 """Return the set of flags appropriate for a quantum for which
129 `PipelineTaskConnections.adjustdQuantum` raised `NoWorkFound`.
130 """
131 return cls.NO_WORK | cls.ADJUST_QUANTUM_RAISED | cls.ANY_OUTPUTS_MISSING | cls.ALL_OUTPUTS_MISSING
133 def concise(self) -> str:
134 """Return a concise string representation of the flags.
136 Returns
137 -------
138 s : `str`
139 Two-character string representation, with the first character
140 indicating whether any predicted outputs were missing and the
141 second representing any exceptions raised. This representation is
142 not always complete; some rare combinations of flags are displayed
143 as if only one of the flags was set.
145 Notes
146 -----
147 The `legend` method returns a description of the returned codes.
148 """
149 char1 = ""
150 if self & QuantumSuccessCaveats.ALL_OUTPUTS_MISSING:
151 char1 = "*"
152 elif self & QuantumSuccessCaveats.ANY_OUTPUTS_MISSING:
153 char1 = "+"
154 char2 = ""
155 if self & QuantumSuccessCaveats.ADJUST_QUANTUM_RAISED:
156 char2 = "A"
157 elif self & QuantumSuccessCaveats.UNPROCESSABLE_DATA:
158 char2 = "D"
159 elif self & QuantumSuccessCaveats.UPSTREAM_FAILURE_NO_WORK:
160 char2 = "U"
161 elif self & QuantumSuccessCaveats.PARTIAL_OUTPUTS_ERROR:
162 char2 = "P"
163 elif self & QuantumSuccessCaveats.NO_WORK: 163 ↛ 165line 163 didn't jump to line 165 because the condition on line 163 was always true
164 char2 = "N"
165 return char1 + char2
167 @staticmethod
168 def legend() -> dict[str, str]:
169 """Return a `dict` with human-readable descriptions of the characters
170 used in `concise`.
172 Returns
173 -------
174 legend : `dict` [ `str`, `str` ]
175 Mapping from character code to description.
176 """
177 return {
178 "+": "at least one predicted output was missing, but not all were",
179 "*": "all predicted outputs were missing (besides logs and metadata)",
180 "A": "adjustQuantum raised NoWorkFound; a regenerated QG would not include this quantum",
181 "D": "algorithm considers data too bad to be processable",
182 "U": "one or more input dataset was incomplete due to an upstream failure",
183 "P": "task failed but wrote partial outputs; considered a partial success",
184 "N": "runQuantum raised NoWorkFound",
185 }
188class ExceptionInfo(pydantic.BaseModel):
189 """Information about an exception that was raised."""
191 type_name: str
192 """Fully-qualified Python type name for the exception raised."""
194 message: str
195 """String message included in the exception."""
197 metadata: dict[str, float | int | str | bool | None]
198 """Additional metadata included in the exception."""
200 @classmethod
201 def _from_metadata(cls, md: TaskMetadata) -> ExceptionInfo:
202 """Construct from task metadata.
204 Parameters
205 ----------
206 md : `TaskMetadata`
207 Metadata about the error, as written by
208 `AnnotatedPartialOutputsError`.
210 Returns
211 -------
212 info : `ExceptionInfo`
213 Information about the exception.
214 """
215 result = cls(type_name=md["type"], message=md["message"], metadata={})
216 if "metadata" in md: 216 ↛ 228line 216 didn't jump to line 228 because the condition on line 216 was always true
217 raw_err_metadata = md["metadata"].to_dict()
218 for k, v in raw_err_metadata.items():
219 # Guard against error metadata we wouldn't be able to serialize
220 # later via Pydantic; don't want one weird value bringing down
221 # our ability to report on an entire run.
222 if isinstance(v, float | int | str | bool): 222 ↛ 225line 222 didn't jump to line 225 because the condition on line 222 was always true
223 result.metadata[k] = v
224 else:
225 _LOG.debug(
226 "Not propagating nested or JSON-incompatible exception metadata key %s=%r.", k, v
227 )
228 return result
230 # Work around the fact that Sphinx chokes on Pydantic docstring formatting,
231 # when we inherit those docstrings in our public classes.
232 if "sphinx" in sys.modules and not TYPE_CHECKING:
234 def copy(self, *args: Any, **kwargs: Any) -> Any:
235 """See `pydantic.BaseModel.copy`."""
236 return super().copy(*args, **kwargs)
238 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
239 """See `pydantic.BaseModel.model_dump`."""
240 return super().model_dump(*args, **kwargs)
242 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any:
243 """See `pydantic.BaseModel.model_dump_json`."""
244 return super().model_dump(*args, **kwargs)
246 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
247 """See `pydantic.BaseModel.model_copy`."""
248 return super().model_copy(*args, **kwargs)
250 @classmethod
251 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override]
252 """See `pydantic.BaseModel.model_construct`."""
253 return super().model_construct(*args, **kwargs)
255 @classmethod
256 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any:
257 """See `pydantic.BaseModel.model_json_schema`."""
258 return super().model_json_schema(*args, **kwargs)
260 @classmethod
261 def model_validate(cls, *args: Any, **kwargs: Any) -> Any:
262 """See `pydantic.BaseModel.model_validate`."""
263 return super().model_validate(*args, **kwargs)
265 @classmethod
266 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any:
267 """See `pydantic.BaseModel.model_validate_json`."""
268 return super().model_validate_json(*args, **kwargs)
270 @classmethod
271 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any:
272 """See `pydantic.BaseModel.model_validate_strings`."""
273 return super().model_validate_strings(*args, **kwargs)
276class QuantumAttemptStatus(enum.Enum):
277 """Enum summarizing an attempt to run a quantum."""
279 ABORTED = -4
280 """The quantum failed with a hard error that prevented both logs and
281 metadata from being written.
283 This state is only set if information from higher-level tooling (e.g. BPS)
284 is available to distinguish it from ``UNKNOWN``.
285 """
287 UNKNOWN = -3
288 """The status of this attempt is unknown.
290 This means no logs or metadata were written, and it at least could not be
291 determined whether the quantum was blocked by an upstream failure (if it
292 was definitely blocked, `BLOCKED` is set instead).
293 """
295 ABORTED_SUCCESS = -2
296 """Task metadata was written for this attempt but logs were not.
298 This is a rare condition that requires a hard failure (i.e. the kind that
299 can prevent a ``finally`` block from running or I/O from being durable) at
300 a very precise time.
301 """
303 FAILED = -1
304 """Execution of the quantum failed gracefully.
306 This is always set if the task metadata dataset was not written but logs
307 were, as is the case when a Python exception is caught and handled by the
308 execution system.
310 This status guarantees that the task log dataset was produced but the
311 metadata dataset was not.
312 """
314 BLOCKED = 0
315 """This quantum was not executed because an upstream quantum failed.
317 Upstream quanta with status `UNKNOWN`, `FAILED`, or `ABORTED` are
318 considered blockers; `ABORTED_SUCCESS` is not.
319 """
321 SUCCESSFUL = 1
322 """This quantum was successfully executed.
324 Quanta may be considered successful even if they do not write any outputs
325 or shortcut early by raising `NoWorkFound` or one of its variants. They
326 may even be considered successful if they raise
327 `AnnotatedPartialOutputsError` if the executor is configured to treat that
328 exception as a non-failure. See `QuantumSuccessCaveats` for details on how
329 these "successes with caveats" are reported.
330 """
332 @property
333 def has_metadata(self) -> bool:
334 """Whether the task metadata dataset was produced."""
335 return self is self.SUCCESSFUL or self is self.ABORTED_SUCCESS
337 @property
338 def has_log(self) -> bool:
339 """Whether the log dataset was produced."""
340 return self is self.SUCCESSFUL or self is self.FAILED
342 @property
343 def title(self) -> str:
344 """A version of this status' name suitable for use as a title in a plot
345 or table.
346 """
347 return self.name.capitalize().replace("_", " ")
349 @property
350 def is_rare(self) -> bool:
351 """Whether this status is rare enough that it should only be listed
352 when it actually occurs.
353 """
354 return self in (self.ABORTED, self.ABORTED_SUCCESS, self.UNKNOWN)
357class GetSetDictMetadataHolder(Protocol):
358 """Protocol for objects that have a ``metadata`` attribute that satisfies
359 `GetSetDictMetadata`.
360 """
362 @property
363 def metadata(self) -> GetSetDictMetadata | MutableMapping[str, Any] | None:
364 pass
367class NoWorkFound(BaseException):
368 """An exception raised when a Quantum should not exist because there is no
369 work for it to do.
371 This usually occurs because a non-optional input dataset is not present, or
372 a spatiotemporal overlap that was conservatively predicted does not
373 actually exist.
375 This inherits from BaseException because it is used to signal a case that
376 we don't consider a real error, even though we often want to use try/except
377 logic to trap it.
378 """
380 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK
383class UpstreamFailureNoWorkFound(NoWorkFound):
384 """A specialization of `NoWorkFound` that indicates that an upstream task
385 had a problem that was ignored (e.g. to prevent a single-detector failure
386 from bringing down an entire visit).
387 """
389 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK | QuantumSuccessCaveats.UPSTREAM_FAILURE_NO_WORK
392class RepeatableQuantumError(RuntimeError):
393 """Exception that may be raised by PipelineTasks (and code they delegate
394 to) in order to indicate that a repeatable problem that will not be
395 addressed by retries.
397 This usually indicates that the algorithm and the data it has been given
398 are somehow incompatible, and the task should run fine on most other data.
400 This exception may be used as a base class for more specific questions, or
401 used directly while chaining another exception, e.g.::
403 try:
404 run_code()
405 except SomeOtherError as err:
406 raise RepeatableQuantumError() from err
408 This may be used for missing input data when the desired behavior is to
409 cause all downstream tasks being run be blocked, forcing the user to
410 address the problem. When the desired behavior is to skip all of this
411 quantum and attempt downstream tasks (or skip them) without its its
412 outputs, raise `NoWorkFound` or return without raising instead.
413 """
415 EXIT_CODE = 20
418class AlgorithmError(RepeatableQuantumError, abc.ABC):
419 """Exception that may be raised by PipelineTasks (and code they delegate
420 to) in order to indicate a repeatable algorithmic failure that will not be
421 addressed by retries.
423 Subclass this exception to define the metadata associated with the error
424 (for example: number of data points in a fit vs. degrees of freedom).
425 """
427 def __new__(cls, *args: Any, **kwargs: Any) -> AlgorithmError:
428 # Have to override __new__ because builtin subclasses aren't checked
429 # for abstract methods; see https://github.com/python/cpython/issues/50246
430 if cls.__abstractmethods__:
431 raise TypeError(
432 f"Can't instantiate abstract class {cls.__name__} with "
433 f"abstract methods: {','.join(sorted(cls.__abstractmethods__))}"
434 )
435 return super().__new__(cls, *args, **kwargs)
437 @property
438 @abc.abstractmethod
439 def metadata(self) -> NestedMetadataDict | None:
440 """Metadata from the raising `~lsst.pipe.base.Task` with more
441 information about the failure. The contents of the dict are
442 `~lsst.pipe.base.Task`-dependent, and must have `str` keys and `str`,
443 `int`, `float`, `bool`, or nested-dictionary (with the same key and
444 value types) values.
445 """
446 raise NotImplementedError
449class UnprocessableDataError(NoWorkFound):
450 """A specialization of `NoWorkFound` that will be [subclassed and] raised
451 by Tasks to indicate a failure to process their inputs for some reason that
452 is non-recoverable.
454 Notes
455 -----
456 An example is a known bright star that causes PSF measurement to fail, and
457 that makes that detector entirely non-recoverable. Another example is an
458 image with an oddly shaped PSF (e.g. due to a failure to achieve focus)
459 that warrants the image being flagged as "poor quality" which should not
460 have further processing attempted.
462 The `NoWorkFound` inheritance ensures the job will not be considered a
463 failure (i.e. such that no human time will inadvertently be spent chasing
464 it down).
466 Do not raise this unless we are convinced that the data cannot (or should
467 not) be processed, even by a better algorithm. Most instances where this
468 error would be raised likely require an RFC to explicitly define the
469 situation.
470 """
472 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK | QuantumSuccessCaveats.UNPROCESSABLE_DATA
475class AnnotatedPartialOutputsError(RepeatableQuantumError):
476 """Exception that runQuantum raises when the (partial) outputs it has
477 written contain information about their own incompleteness or degraded
478 quality.
480 Clients should construct this exception by calling `annotate` instead of
481 calling the constructor directly. However, `annotate` does not chain the
482 exception; this must still be done by the client.
484 This exception should always chain the original error. When the
485 executor catches this exception, it will report the original exception. In
486 contrast, other exceptions raised from ``runQuantum`` are considered to
487 invalidate any outputs that are already written.
488 """
490 FLAGS: ClassVar = QuantumSuccessCaveats.PARTIAL_OUTPUTS_ERROR
492 @classmethod
493 def annotate(
494 cls, error: Exception, *args: GetSetDictMetadataHolder | None, log: logging.Logger | LsstLogAdapter
495 ) -> AnnotatedPartialOutputsError:
496 """Set metadata on outputs to explain the nature of the failure.
498 Parameters
499 ----------
500 error : `Exception`
501 Exception that caused the task to fail.
502 *args : `GetSetDictMetadataHolder`
503 Objects (e.g. Task, Exposure, SimpleCatalog) to annotate with
504 failure information. They must have a `metadata` property that
505 is a `~collections.abc.MutableMapping`.
506 log : `logging.Logger`
507 Log to send error message to.
509 Returns
510 -------
511 error : `AnnotatedPartialOutputsError`
512 Exception that the failing task can ``raise from`` with the
513 passed-in exception.
515 Notes
516 -----
517 This should be called from within an except block that has caught an
518 exception. Here is an example of handling a failure in
519 ``PipelineTask.runQuantum`` that annotates and writes partial outputs:
521 .. code-block:: py
522 :name: annotate-error-example
524 def runQuantum(self, butlerQC, inputRefs, outputRefs):
525 inputs = butlerQC.get(inputRefs)
526 exposures = inputs.pop("exposures")
527 assert not inputs, "runQuantum got more inputs than expected"
529 result = pipeBase.Struct(catalog=None)
530 try:
531 self.run(exposure)
532 except pipeBase.AlgorithmError as e:
533 error = pipeBase.AnnotatedPartialOutputsError.annotate(
534 e, self, result.catalog, log=self.log
535 )
536 raise error from e
537 finally:
538 butlerQC.put(result, outputRefs)
539 """
540 failure_info = {
541 "message": str(error),
542 "type": introspection.get_full_type_name(error),
543 }
544 if other := getattr(error, "metadata", None):
545 failure_info["metadata"] = other
547 # NOTE: Can't fully test this in pipe_base because afw is not a
548 # dependency; test_calibrateImage.py in pipe_tasks gives more coverage.
549 for item in args:
550 # Some outputs may not exist, so we cannot set metadata on them.
551 if item is None: 551 ↛ 552line 551 didn't jump to line 552 because the condition on line 551 was never true
552 continue
553 if hasattr(item.metadata, "set_dict"): 553 ↛ 555line 553 didn't jump to line 555 because the condition on line 553 was always true
554 item.metadata.set_dict("failure", failure_info) # type: ignore
555 elif isinstance(item.metadata, MutableMapping):
556 item.metadata["failure"] = failure_info
557 else:
558 raise TypeError(f"Invalid metadata type {type(item.metadata).__name__!r}")
560 log.debug(
561 "Task failed with only partial outputs; see exception message for details.",
562 exc_info=error,
563 )
565 return cls("Task failed and wrote partial outputs: see chained exception for details.")
568class InvalidQuantumError(Exception):
569 """Exception that may be raised by PipelineTasks (and code they delegate
570 to) in order to indicate logic bug or configuration problem.
572 This usually indicates that the configured algorithm itself is invalid and
573 will not run on a significant fraction of quanta (often all of them).
575 This exception may be used as a base class for more specific questions, or
576 used directly while chaining another exception, e.g.::
578 try:
579 run_code()
580 except SomeOtherError as err:
581 raise RepeatableQuantumError() from err
583 Raising this exception in `PipelineTask.runQuantum` or something it calls
584 is a last resort - whenever possible, such problems should cause exceptions
585 in ``__init__`` or in QuantumGraph generation. It should never be used
586 for missing data.
587 """
589 EXIT_CODE = 21