Coverage for python/lsst/pipe/base/_status.py: 92%

141 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-22 09:37 +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/>. 

27 

28from __future__ import annotations 

29 

30__all__ = ( 

31 "AlgorithmError", 

32 "AnnotatedPartialOutputsError", 

33 "ExceptionInfo", 

34 "InvalidQuantumError", 

35 "NoWorkFound", 

36 "QuantumAttemptStatus", 

37 "QuantumSuccessCaveats", 

38 "RepeatableQuantumError", 

39 "UnprocessableDataError", 

40 "UpstreamFailureNoWorkFound", 

41) 

42 

43import abc 

44import enum 

45import logging 

46import re 

47import sys 

48from collections.abc import MutableMapping 

49from typing import TYPE_CHECKING, Any, ClassVar, Protocol 

50 

51import pydantic 

52 

53from lsst.utils import introspection 

54from lsst.utils.logging import LsstLogAdapter, getLogger 

55 

56from ._task_metadata import GetSetDictMetadata, NestedMetadataDict 

57 

58if TYPE_CHECKING: 

59 from ._task_metadata import TaskMetadata 

60 

61 

62_LOG = getLogger(__name__) 

63 

64 

65class QuantumSuccessCaveats(enum.Flag): 

66 """Flags that add caveats to a "successful" quantum. 

67 

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 """ 

73 

74 NO_CAVEATS = 0 

75 """All outputs were produced and no exceptions were raised.""" 

76 

77 ANY_OUTPUTS_MISSING = enum.auto() 

78 """At least one predicted output was not produced.""" 

79 

80 ALL_OUTPUTS_MISSING = enum.auto() 

81 """No predicted outputs (except logs and metadata) were produced. 

82 

83 `ANY_OUTPUTS_MISSING` is also set whenever this flag is set. 

84 """ 

85 

86 NO_WORK = enum.auto() 

87 """A subclass of `NoWorkFound` was raised. 

88 

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 """ 

93 

94 ADJUST_QUANTUM_RAISED = enum.auto() 

95 """`NoWorkFound` was raised by `PipelineTaskConnnections.adjustQuantum`. 

96 

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. 

101 

102 `NO_WORK` and `ALL_OUTPUTS_MISSING` are also set whenever this flag is set. 

103 """ 

104 

105 UPSTREAM_FAILURE_NO_WORK = enum.auto() 

106 """`UpstreamFailureNoWorkFound` was raised by `PipelineTask.runQuantum`. 

107 

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`. 

111 

112 `NO_WORK` is also set whenever this flag is set. 

113 """ 

114 

115 UNPROCESSABLE_DATA = enum.auto() 

116 """`UnprocessableDataError` was raised by `PipelineTask.runQuantum`. 

117 

118 `NO_WORK` is also set whenever this flag is set. 

119 """ 

120 

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 """ 

126 

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 

133 

134 @classmethod 

135 def expanded_dict(cls, concise: str) -> dict | None: 

136 """Return a dictionary representation of the concise flags. 

137 

138 Parameters 

139 ---------- 

140 concise : `str` 

141 The concise string representation of the flags to expand to a 

142 dictionary. 

143 

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 

154 

155 def concise(self) -> str: 

156 """Return a concise string representation of the flags. 

157 

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. 

166 

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 

188 

189 @staticmethod 

190 def legend() -> dict[str, str]: 

191 """Return a `dict` with human-readable descriptions of the characters 

192 used in `concise`. 

193 

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 } 

208 

209 

210class ExceptionInfo(pydantic.BaseModel): 

211 """Information about an exception that was raised.""" 

212 

213 type_name: str 

214 """Fully-qualified Python type name for the exception raised.""" 

215 

216 message: str 

217 """String message included in the exception.""" 

218 

219 metadata: dict[str, float | int | str | bool | None] 

220 """Additional metadata included in the exception.""" 

221 

222 @classmethod 

223 def _from_metadata(cls, md: TaskMetadata) -> ExceptionInfo: 

224 """Construct from task metadata. 

225 

226 Parameters 

227 ---------- 

228 md : `TaskMetadata` 

229 Metadata about the error, as written by 

230 `AnnotatedPartialOutputsError`. 

231 

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 

251 

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: 

255 

256 def copy(self, *args: Any, **kwargs: Any) -> Any: 

257 """See `pydantic.BaseModel.copy`.""" 

258 return super().copy(*args, **kwargs) 

259 

260 def model_dump(self, *args: Any, **kwargs: Any) -> Any: 

261 """See `pydantic.BaseModel.model_dump`.""" 

262 return super().model_dump(*args, **kwargs) 

263 

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) 

267 

268 def model_copy(self, *args: Any, **kwargs: Any) -> Any: 

269 """See `pydantic.BaseModel.model_copy`.""" 

270 return super().model_copy(*args, **kwargs) 

271 

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) 

276 

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) 

281 

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) 

286 

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) 

291 

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) 

296 

297 

298class QuantumAttemptStatus(enum.Enum): 

299 """Enum summarizing an attempt to run a quantum.""" 

300 

301 ABORTED = -4 

302 """The quantum failed with a hard error that prevented both logs and 

303 metadata from being written. 

304 

305 This state is only set if information from higher-level tooling (e.g. BPS) 

306 is available to distinguish it from ``UNKNOWN``. 

307 """ 

308 

309 UNKNOWN = -3 

310 """The status of this attempt is unknown. 

311 

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 """ 

316 

317 ABORTED_SUCCESS = -2 

318 """Task metadata was written for this attempt but logs were not. 

319 

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 """ 

324 

325 FAILED = -1 

326 """Execution of the quantum failed gracefully. 

327 

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. 

331 

332 This status guarantees that the task log dataset was produced but the 

333 metadata dataset was not. 

334 """ 

335 

336 BLOCKED = 0 

337 """This quantum was not executed because an upstream quantum failed. 

338 

339 Upstream quanta with status `UNKNOWN`, `FAILED`, or `ABORTED` are 

340 considered blockers; `ABORTED_SUCCESS` is not. 

341 """ 

342 

343 SUCCESSFUL = 1 

344 """This quantum was successfully executed. 

345 

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 """ 

353 

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 

358 

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 

363 

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("_", " ") 

370 

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) 

377 

378 

379class GetSetDictMetadataHolder(Protocol): 

380 """Protocol for objects that have a ``metadata`` attribute that satisfies 

381 `GetSetDictMetadata`. 

382 """ 

383 

384 @property 

385 def metadata(self) -> GetSetDictMetadata | MutableMapping[str, Any] | None: 

386 pass 

387 

388 

389class NoWorkFound(BaseException): 

390 """An exception raised when a Quantum should not exist because there is no 

391 work for it to do. 

392 

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. 

396 

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 """ 

401 

402 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK 

403 

404 

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 """ 

410 

411 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK | QuantumSuccessCaveats.UPSTREAM_FAILURE_NO_WORK 

412 

413 

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. 

418 

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. 

421 

422 This exception may be used as a base class for more specific questions, or 

423 used directly while chaining another exception, e.g.:: 

424 

425 try: 

426 run_code() 

427 except SomeOtherError as err: 

428 raise RepeatableQuantumError() from err 

429 

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 """ 

436 

437 EXIT_CODE = 20 

438 

439 

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. 

444 

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 """ 

448 

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) 

458 

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 

469 

470 

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. 

475 

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. 

483 

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). 

487 

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 """ 

493 

494 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK | QuantumSuccessCaveats.UNPROCESSABLE_DATA 

495 

496 

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. 

501 

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. 

505 

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 """ 

511 

512 FLAGS: ClassVar = QuantumSuccessCaveats.PARTIAL_OUTPUTS_ERROR 

513 

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. 

519 

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. 

530 

531 Returns 

532 ------- 

533 error : `AnnotatedPartialOutputsError` 

534 Exception that the failing task can ``raise from`` with the 

535 passed-in exception. 

536 

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: 

542 

543 .. code-block:: py 

544 :name: annotate-error-example 

545 

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" 

550 

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 

568 

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}") 

581 

582 log.debug( 

583 "Task failed with only partial outputs; see exception message for details.", 

584 exc_info=error, 

585 ) 

586 

587 return cls("Task failed and wrote partial outputs: see chained exception for details.") 

588 

589 

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. 

593 

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). 

596 

597 This exception may be used as a base class for more specific questions, or 

598 used directly while chaining another exception, e.g.:: 

599 

600 try: 

601 run_code() 

602 except SomeOtherError as err: 

603 raise RepeatableQuantumError() from err 

604 

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 """ 

610 

611 EXIT_CODE = 21