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:35 +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 sys 

47from collections.abc import MutableMapping 

48from typing import TYPE_CHECKING, Any, ClassVar, Protocol 

49 

50import pydantic 

51 

52from lsst.utils import introspection 

53from lsst.utils.logging import LsstLogAdapter, getLogger 

54 

55from ._task_metadata import GetSetDictMetadata, NestedMetadataDict 

56 

57if TYPE_CHECKING: 

58 from ._task_metadata import TaskMetadata 

59 

60 

61_LOG = getLogger(__name__) 

62 

63 

64class QuantumSuccessCaveats(enum.Flag): 

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

66 

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

72 

73 NO_CAVEATS = 0 

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

75 

76 ANY_OUTPUTS_MISSING = enum.auto() 

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

78 

79 ALL_OUTPUTS_MISSING = enum.auto() 

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

81 

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

83 """ 

84 

85 NO_WORK = enum.auto() 

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

87 

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

92 

93 ADJUST_QUANTUM_RAISED = enum.auto() 

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

95 

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. 

100 

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

102 """ 

103 

104 UPSTREAM_FAILURE_NO_WORK = enum.auto() 

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

106 

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

110 

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

112 """ 

113 

114 UNPROCESSABLE_DATA = enum.auto() 

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

116 

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

118 """ 

119 

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

125 

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 

132 

133 def concise(self) -> str: 

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

135 

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. 

144 

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 

166 

167 @staticmethod 

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

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

170 used in `concise`. 

171 

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 } 

186 

187 

188class ExceptionInfo(pydantic.BaseModel): 

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

190 

191 type_name: str 

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

193 

194 message: str 

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

196 

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

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

199 

200 @classmethod 

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

202 """Construct from task metadata. 

203 

204 Parameters 

205 ---------- 

206 md : `TaskMetadata` 

207 Metadata about the error, as written by 

208 `AnnotatedPartialOutputsError`. 

209 

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 

229 

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: 

233 

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

235 """See `pydantic.BaseModel.copy`.""" 

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

237 

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

239 """See `pydantic.BaseModel.model_dump`.""" 

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

241 

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) 

245 

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

247 """See `pydantic.BaseModel.model_copy`.""" 

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

249 

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) 

254 

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) 

259 

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) 

264 

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) 

269 

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) 

274 

275 

276class QuantumAttemptStatus(enum.Enum): 

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

278 

279 ABORTED = -4 

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

281 metadata from being written. 

282 

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

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

285 """ 

286 

287 UNKNOWN = -3 

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

289 

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

294 

295 ABORTED_SUCCESS = -2 

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

297 

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

302 

303 FAILED = -1 

304 """Execution of the quantum failed gracefully. 

305 

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. 

309 

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

311 metadata dataset was not. 

312 """ 

313 

314 BLOCKED = 0 

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

316 

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

318 considered blockers; `ABORTED_SUCCESS` is not. 

319 """ 

320 

321 SUCCESSFUL = 1 

322 """This quantum was successfully executed. 

323 

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

331 

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 

336 

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 

341 

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

348 

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) 

355 

356 

357class GetSetDictMetadataHolder(Protocol): 

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

359 `GetSetDictMetadata`. 

360 """ 

361 

362 @property 

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

364 pass 

365 

366 

367class NoWorkFound(BaseException): 

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

369 work for it to do. 

370 

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. 

374 

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

379 

380 FLAGS: ClassVar = QuantumSuccessCaveats.NO_WORK 

381 

382 

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

388 

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

390 

391 

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. 

396 

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. 

399 

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

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

402 

403 try: 

404 run_code() 

405 except SomeOtherError as err: 

406 raise RepeatableQuantumError() from err 

407 

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

414 

415 EXIT_CODE = 20 

416 

417 

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. 

422 

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

426 

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) 

436 

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 

447 

448 

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. 

453 

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. 

461 

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

465 

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

471 

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

473 

474 

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. 

479 

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. 

483 

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

489 

490 FLAGS: ClassVar = QuantumSuccessCaveats.PARTIAL_OUTPUTS_ERROR 

491 

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. 

497 

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. 

508 

509 Returns 

510 ------- 

511 error : `AnnotatedPartialOutputsError` 

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

513 passed-in exception. 

514 

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: 

520 

521 .. code-block:: py 

522 :name: annotate-error-example 

523 

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" 

528 

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 

546 

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

559 

560 log.debug( 

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

562 exc_info=error, 

563 ) 

564 

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

566 

567 

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. 

571 

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

574 

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

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

577 

578 try: 

579 run_code() 

580 except SomeOtherError as err: 

581 raise RepeatableQuantumError() from err 

582 

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

588 

589 EXIT_CODE = 21