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

729 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-29 02:11 -0700

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

32 "ProvenanceDatasetModel", 

33 "ProvenanceInitQuantumInfo", 

34 "ProvenanceInitQuantumModel", 

35 "ProvenanceLogRecordsModel", 

36 "ProvenanceQuantumGraph", 

37 "ProvenanceQuantumGraphReader", 

38 "ProvenanceQuantumGraphWriter", 

39 "ProvenanceQuantumInfo", 

40 "ProvenanceQuantumModel", 

41 "ProvenanceQuantumReport", 

42 "ProvenanceQuantumScanData", 

43 "ProvenanceQuantumScanModels", 

44 "ProvenanceQuantumScanStatus", 

45 "ProvenanceReport", 

46 "ProvenanceTaskMetadataModel", 

47) 

48 

49import dataclasses 

50import enum 

51import itertools 

52import json 

53import sys 

54import uuid 

55from collections import Counter 

56from collections.abc import Callable, Iterable, Iterator, Mapping 

57from contextlib import ExitStack, contextmanager 

58from typing import TYPE_CHECKING, Any, Literal, TypedDict 

59 

60import astropy.table 

61import networkx 

62import numpy as np 

63import pydantic 

64 

65from lsst.daf.butler import Butler, DataCoordinate 

66from lsst.daf.butler.logging import ButlerLogRecord, ButlerLogRecords 

67from lsst.resources import ResourcePath, ResourcePathExpression 

68from lsst.utils.iteration import ensure_iterable 

69from lsst.utils.logging import LsstLogAdapter, getLogger 

70from lsst.utils.packages import Packages 

71 

72from .. import automatic_connection_constants as acc 

73from .._status import ExceptionInfo, QuantumAttemptStatus, QuantumSuccessCaveats 

74from .._task_metadata import TaskMetadata 

75from ..log_capture import _ExecutionLogRecordsExtra 

76from ..log_on_close import LogOnClose 

77from ..pipeline_graph import PipelineGraph, TaskImportMode, TaskInitNode 

78from ..resource_usage import QuantumResourceUsage 

79from ._common import ( 

80 BaseQuantumGraph, 

81 BaseQuantumGraphReader, 

82 BaseQuantumGraphWriter, 

83 ConnectionName, 

84 DataCoordinateValues, 

85 DatasetInfo, 

86 DatasetTypeName, 

87 HeaderModel, 

88 QuantumInfo, 

89 TaskLabel, 

90) 

91from ._multiblock import Compressor, MultiblockReader, MultiblockWriter 

92from ._predicted import ( 

93 PredictedDatasetModel, 

94 PredictedQuantumDatasetsModel, 

95 PredictedQuantumGraph, 

96 PredictedQuantumGraphComponents, 

97) 

98 

99# Sphinx needs imports for type annotations of base class members. 

100if "sphinx" in sys.modules: 

101 import zipfile # noqa: F401 

102 

103 from ._multiblock import AddressReader, Decompressor # noqa: F401 

104 

105 

106type LoopWrapper[T] = Callable[[Iterable[T]], Iterable[T]] 

107 

108_LOG = getLogger(__file__) 

109 

110DATASET_ADDRESS_INDEX = 0 

111QUANTUM_ADDRESS_INDEX = 1 

112LOG_ADDRESS_INDEX = 2 

113METADATA_ADDRESS_INDEX = 3 

114 

115DATASET_MB_NAME = "datasets" 

116QUANTUM_MB_NAME = "quanta" 

117LOG_MB_NAME = "logs" 

118METADATA_MB_NAME = "metadata" 

119 

120 

121def pass_through[T](arg: T) -> T: 

122 return arg 

123 

124 

125class ProvenanceDatasetInfo(DatasetInfo): 

126 """A typed dictionary that annotates the attributes of the NetworkX graph 

127 node data for a provenance dataset. 

128 

129 Since NetworkX types are not generic over their node mapping type, this has 

130 to be used explicitly, e.g.:: 

131 

132 node_data: ProvenanceDatasetInfo = xgraph.nodes[dataset_id] 

133 

134 where ``xgraph`` is `ProvenanceQuantumGraph.bipartite_xgraph`. 

135 """ 

136 

137 dataset_id: uuid.UUID 

138 """Unique identifier for the dataset.""" 

139 

140 produced: bool 

141 """Whether this dataset was produced (vs. only predicted). 

142 

143 This is always `True` for overall input datasets. It is also `True` for 

144 datasets that were produced and then removed before/during transfer back to 

145 the central butler repository, so it may not reflect the continued 

146 existence of the dataset. 

147 """ 

148 

149 

150class ProvenanceQuantumInfo(QuantumInfo): 

151 """A typed dictionary that annotates the attributes of the NetworkX graph 

152 node data for a provenance quantum. 

153 

154 Since NetworkX types are not generic over their node mapping type, this has 

155 to be used explicitly, e.g.:: 

156 

157 node_data: ProvenanceQuantumInfo = xgraph.nodes[quantum_id] 

158 

159 where ``xgraph`` is `ProvenanceQuantumGraph.bipartite_xgraph` or 

160 `ProvenanceQuantumGraph.quantum_only_xgraph` 

161 """ 

162 

163 status: QuantumAttemptStatus 

164 """Enumerated status for the quantum. 

165 

166 This corresponds to the last attempt to run this quantum, or 

167 `QuantumAttemptStatus.BLOCKED` if there were no attempts. 

168 """ 

169 

170 caveats: QuantumSuccessCaveats | None 

171 """Flags indicating caveats on successful quanta. 

172 

173 This corresponds to the last attempt to run this quantum. 

174 """ 

175 

176 exception: ExceptionInfo | None 

177 """Information about an exception raised when the quantum was executing. 

178 

179 This corresponds to the last attempt to run this quantum. 

180 """ 

181 

182 resource_usage: QuantumResourceUsage | None 

183 """Resource usage information (timing, memory use) for this quantum. 

184 

185 This corresponds to the last attempt to run this quantum. 

186 """ 

187 

188 attempts: list[ProvenanceQuantumAttemptModel] 

189 """Information about each attempt to run this quantum. 

190 

191 An entry is added merely if the quantum *should* have been attempted; an 

192 empty `list` is used only for quanta that were blocked by an upstream 

193 failure. 

194 """ 

195 

196 metadata_id: uuid.UUID 

197 """ID of this quantum's metadata dataset.""" 

198 

199 log_id: uuid.UUID 

200 """ID of this quantum's log dataset.""" 

201 

202 

203class ProvenanceInitQuantumInfo(TypedDict): 

204 """A typed dictionary that annotates the attributes of the NetworkX graph 

205 node data for a provenance init quantum. 

206 

207 Since NetworkX types are not generic over their node mapping type, this has 

208 to be used explicitly, e.g.:: 

209 

210 node_data: ProvenanceInitQuantumInfo = xgraph.nodes[quantum_id] 

211 

212 where ``xgraph`` is `ProvenanceQuantumGraph.bipartite_xgraph`. 

213 """ 

214 

215 data_id: DataCoordinate 

216 """Data ID of the quantum. 

217 

218 This is always an empty ID; this key exists to allow init-quanta and 

219 regular quanta to be treated more similarly. 

220 """ 

221 

222 task_label: str 

223 """Label of the task for this quantum.""" 

224 

225 pipeline_node: TaskInitNode 

226 """Node in the pipeline graph for this task's init-only step.""" 

227 

228 config_id: uuid.UUID 

229 """ID of this task's config dataset.""" 

230 

231 

232class ProvenanceDatasetModel(PredictedDatasetModel): 

233 """Data model for the datasets in a provenance quantum graph file.""" 

234 

235 produced: bool 

236 """Whether this dataset was produced (vs. only predicted). 

237 

238 This is always `True` for overall input datasets. It is also `True` for 

239 datasets that were produced and then removed before/during transfer back to 

240 the central butler repository, so it may not reflect the continued 

241 existence of the dataset. 

242 """ 

243 

244 producer: uuid.UUID | None = None 

245 """ID of the quantum that produced this dataset. 

246 

247 This is `None` for overall inputs to the graph. 

248 """ 

249 

250 consumers: list[uuid.UUID] = pydantic.Field(default_factory=list) 

251 """IDs of quanta that were predicted to consume this dataset.""" 

252 

253 @property 

254 def node_id(self) -> uuid.UUID: 

255 """Alias for the dataset ID.""" 

256 return self.dataset_id 

257 

258 @classmethod 

259 def from_predicted( 

260 cls, 

261 predicted: PredictedDatasetModel, 

262 producer: uuid.UUID | None = None, 

263 consumers: Iterable[uuid.UUID] = (), 

264 ) -> ProvenanceDatasetModel: 

265 """Construct from a predicted dataset model. 

266 

267 Parameters 

268 ---------- 

269 predicted : `PredictedDatasetModel` 

270 Information about the dataset from the predicted graph. 

271 producer : `uuid.UUID` or `None`, optional 

272 ID of the quantum that was predicted to produce this dataset. 

273 consumers : `~collections.abc.Iterable` [`uuid.UUID`], optional 

274 IDs of the quanta that were predicted to consume this dataset. 

275 

276 Returns 

277 ------- 

278 provenance : `ProvenanceDatasetModel` 

279 Provenance dataset model. 

280 

281 Notes 

282 ----- 

283 This initializes `produced` to `True` when ``producer is None`` and 

284 `False` otherwise, on the assumption that it will be updated later. 

285 """ 

286 return cls.model_construct( 

287 dataset_id=predicted.dataset_id, 

288 dataset_type_name=predicted.dataset_type_name, 

289 data_coordinate=predicted.data_coordinate, 

290 run=predicted.run, 

291 produced=(producer is None), # if it's not produced by this QG, it's an overall input 

292 producer=producer, 

293 consumers=list(consumers), 

294 ) 

295 

296 def _add_to_graph(self, graph: ProvenanceQuantumGraph) -> None: 

297 """Add this dataset and its edges to quanta to a provenance graph. 

298 

299 Parameters 

300 ---------- 

301 graph : `ProvenanceQuantumGraph` 

302 Graph to update in place. 

303 

304 Notes 

305 ----- 

306 This method adds: 

307 

308 - a ``bipartite_xgraph`` dataset node with full attributes; 

309 - ``bipartite_xgraph`` edges to adjacent quanta (which adds quantum 

310 nodes with no attributes), without populating edge attributes; 

311 - ``quantum_only_xgraph`` edges for each pair of quanta in which one 

312 produces this dataset and another consumes it (this also adds quantum 

313 nodes with no attributes). 

314 """ 

315 dataset_type_node = graph.pipeline_graph.dataset_types[self.dataset_type_name] 

316 data_id = DataCoordinate.from_full_values(dataset_type_node.dimensions, tuple(self.data_coordinate)) 

317 graph._bipartite_xgraph.add_node( 

318 self.dataset_id, 

319 data_id=data_id, 

320 dataset_type_name=self.dataset_type_name, 

321 pipeline_node=dataset_type_node, 

322 run=self.run, 

323 produced=self.produced, 

324 ) 

325 if self.producer is not None: 

326 graph._bipartite_xgraph.add_edge(self.producer, self.dataset_id) 

327 for consumer_id in self.consumers: 

328 graph._bipartite_xgraph.add_edge(self.dataset_id, consumer_id) 

329 if self.producer is not None: 

330 graph._quantum_only_xgraph.add_edge(self.producer, consumer_id) 

331 graph._datasets_by_type[self.dataset_type_name][data_id] = self.dataset_id 

332 

333 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

334 # when we inherit those docstrings in our public classes. 

335 if "sphinx" in sys.modules and not TYPE_CHECKING: 

336 

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

338 """See `pydantic.BaseModel.copy`.""" 

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

340 

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

342 """See `pydantic.BaseModel.model_dump`.""" 

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

344 

345 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

346 """See `pydantic.BaseModel.model_dump_json`.""" 

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

348 

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

350 """See `pydantic.BaseModel.model_copy`.""" 

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

352 

353 @classmethod 

354 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

355 """See `pydantic.BaseModel.model_construct`.""" 

356 return super().model_construct(*args, **kwargs) 

357 

358 @classmethod 

359 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

360 """See `pydantic.BaseModel.model_json_schema`.""" 

361 return super().model_json_schema(*args, **kwargs) 

362 

363 @classmethod 

364 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

365 """See `pydantic.BaseModel.model_validate`.""" 

366 return super().model_validate(*args, **kwargs) 

367 

368 @classmethod 

369 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

370 """See `pydantic.BaseModel.model_validate_json`.""" 

371 return super().model_validate_json(*args, **kwargs) 

372 

373 @classmethod 

374 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

375 """See `pydantic.BaseModel.model_validate_strings`.""" 

376 return super().model_validate_strings(*args, **kwargs) 

377 

378 

379class ProvenanceQuantumAttemptModel(pydantic.BaseModel): 

380 """Data model for a now-superseded attempt to run a quantum in a 

381 provenance quantum graph file. 

382 """ 

383 

384 attempt: int = 0 

385 """Counter incremented for every attempt to execute this quantum.""" 

386 

387 status: QuantumAttemptStatus = QuantumAttemptStatus.UNKNOWN 

388 """Enumerated status for the quantum.""" 

389 

390 caveats: QuantumSuccessCaveats | None = None 

391 """Flags indicating caveats on successful quanta.""" 

392 

393 exception: ExceptionInfo | None = None 

394 """Information about an exception raised when the quantum was executing.""" 

395 

396 resource_usage: QuantumResourceUsage | None = None 

397 """Resource usage information (timing, memory use) for this quantum.""" 

398 

399 previous_process_quanta: list[uuid.UUID] = pydantic.Field(default_factory=list) 

400 """The IDs of other quanta previously executed in the same process as this 

401 one. 

402 """ 

403 

404 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

405 # when we inherit those docstrings in our public classes. 

406 if "sphinx" in sys.modules and not TYPE_CHECKING: 

407 

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

409 """See `pydantic.BaseModel.copy`.""" 

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

411 

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

413 """See `pydantic.BaseModel.model_dump`.""" 

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

415 

416 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

417 """See `pydantic.BaseModel.model_dump_json`.""" 

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

419 

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

421 """See `pydantic.BaseModel.model_copy`.""" 

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

423 

424 @classmethod 

425 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

426 """See `pydantic.BaseModel.model_construct`.""" 

427 return super().model_construct(*args, **kwargs) 

428 

429 @classmethod 

430 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

431 """See `pydantic.BaseModel.model_json_schema`.""" 

432 return super().model_json_schema(*args, **kwargs) 

433 

434 @classmethod 

435 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

436 """See `pydantic.BaseModel.model_validate`.""" 

437 return super().model_validate(*args, **kwargs) 

438 

439 @classmethod 

440 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

441 """See `pydantic.BaseModel.model_validate_json`.""" 

442 return super().model_validate_json(*args, **kwargs) 

443 

444 @classmethod 

445 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

446 """See `pydantic.BaseModel.model_validate_strings`.""" 

447 return super().model_validate_strings(*args, **kwargs) 

448 

449 

450class ProvenanceLogRecordsModel(pydantic.BaseModel): 

451 """Data model for storing execution logs in a provenance quantum graph 

452 file. 

453 """ 

454 

455 attempts: list[list[ButlerLogRecord] | None] = pydantic.Field(default_factory=list) 

456 """Logs from attempts to run this task, ordered chronologically from first 

457 to last. 

458 """ 

459 

460 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

461 # when we inherit those docstrings in our public classes. 

462 if "sphinx" in sys.modules and not TYPE_CHECKING: 

463 

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

465 """See `pydantic.BaseModel.copy`.""" 

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

467 

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

469 """See `pydantic.BaseModel.model_dump`.""" 

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

471 

472 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

473 """See `pydantic.BaseModel.model_dump_json`.""" 

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

475 

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

477 """See `pydantic.BaseModel.model_copy`.""" 

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

479 

480 @classmethod 

481 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

482 """See `pydantic.BaseModel.model_construct`.""" 

483 return super().model_construct(*args, **kwargs) 

484 

485 @classmethod 

486 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

487 """See `pydantic.BaseModel.model_json_schema`.""" 

488 return super().model_json_schema(*args, **kwargs) 

489 

490 @classmethod 

491 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

492 """See `pydantic.BaseModel.model_validate`.""" 

493 return super().model_validate(*args, **kwargs) 

494 

495 @classmethod 

496 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

497 """See `pydantic.BaseModel.model_validate_json`.""" 

498 return super().model_validate_json(*args, **kwargs) 

499 

500 @classmethod 

501 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

502 """See `pydantic.BaseModel.model_validate_strings`.""" 

503 return super().model_validate_strings(*args, **kwargs) 

504 

505 

506class ProvenanceTaskMetadataModel(pydantic.BaseModel): 

507 """Data model for storing task metadata in a provenance quantum graph 

508 file. 

509 """ 

510 

511 # We want to convert infs and nans to constants, not null. Unfortunately 

512 # the fact that TaskMetadata _also_ sets this is ignored when that model 

513 # is nested here. 

514 model_config = pydantic.ConfigDict(ser_json_inf_nan="constants") 

515 

516 attempts: list[TaskMetadata | None] = pydantic.Field(default_factory=list) 

517 """Metadata from attempts to run this task, ordered chronologically from 

518 first to last. 

519 """ 

520 

521 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

522 # when we inherit those docstrings in our public classes. 

523 if "sphinx" in sys.modules and not TYPE_CHECKING: 

524 

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

526 """See `pydantic.BaseModel.copy`.""" 

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

528 

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

530 """See `pydantic.BaseModel.model_dump`.""" 

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

532 

533 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

534 """See `pydantic.BaseModel.model_dump_json`.""" 

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

536 

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

538 """See `pydantic.BaseModel.model_copy`.""" 

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

540 

541 @classmethod 

542 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

543 """See `pydantic.BaseModel.model_construct`.""" 

544 return super().model_construct(*args, **kwargs) 

545 

546 @classmethod 

547 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

548 """See `pydantic.BaseModel.model_json_schema`.""" 

549 return super().model_json_schema(*args, **kwargs) 

550 

551 @classmethod 

552 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

553 """See `pydantic.BaseModel.model_validate`.""" 

554 return super().model_validate(*args, **kwargs) 

555 

556 @classmethod 

557 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

558 """See `pydantic.BaseModel.model_validate_json`.""" 

559 return super().model_validate_json(*args, **kwargs) 

560 

561 @classmethod 

562 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

563 """See `pydantic.BaseModel.model_validate_strings`.""" 

564 return super().model_validate_strings(*args, **kwargs) 

565 

566 

567class ProvenanceQuantumReport(pydantic.BaseModel): 

568 """A Pydantic model that used to report information about a single 

569 (generally problematic) quantum. 

570 """ 

571 

572 quantum_id: uuid.UUID 

573 data_id: dict[str, int | str] 

574 attempts: list[ProvenanceQuantumAttemptModel] 

575 

576 @classmethod 

577 def from_info(cls, quantum_id: uuid.UUID, quantum_info: ProvenanceQuantumInfo) -> ProvenanceQuantumReport: 

578 """Construct from a provenance quantum graph node. 

579 

580 Parameters 

581 ---------- 

582 quantum_id : `uuid.UUID` 

583 Unique ID for the quantum. 

584 quantum_info : `ProvenanceQuantumInfo` 

585 Node attributes for this quantum. 

586 """ 

587 return cls( 

588 quantum_id=quantum_id, 

589 data_id=dict(quantum_info["data_id"].mapping), 

590 attempts=quantum_info["attempts"], 

591 ) 

592 

593 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

594 # when we inherit those docstrings in our public classes. 

595 if "sphinx" in sys.modules and not TYPE_CHECKING: 

596 

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

598 """See `pydantic.BaseModel.copy`.""" 

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

600 

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

602 """See `pydantic.BaseModel.model_dump`.""" 

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

604 

605 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

606 """See `pydantic.BaseModel.model_dump_json`.""" 

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

608 

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

610 """See `pydantic.BaseModel.model_copy`.""" 

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

612 

613 @classmethod 

614 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

615 """See `pydantic.BaseModel.model_construct`.""" 

616 return super().model_construct(*args, **kwargs) 

617 

618 @classmethod 

619 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

620 """See `pydantic.BaseModel.model_json_schema`.""" 

621 return super().model_json_schema(*args, **kwargs) 

622 

623 @classmethod 

624 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

625 """See `pydantic.BaseModel.model_validate`.""" 

626 return super().model_validate(*args, **kwargs) 

627 

628 @classmethod 

629 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

630 """See `pydantic.BaseModel.model_validate_json`.""" 

631 return super().model_validate_json(*args, **kwargs) 

632 

633 @classmethod 

634 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

635 """See `pydantic.BaseModel.model_validate_strings`.""" 

636 return super().model_validate_strings(*args, **kwargs) 

637 

638 

639class ProvenanceReport(pydantic.RootModel): 

640 """A Pydantic model that groups quantum information by task label, then 

641 status (as a string), and then exception type. 

642 """ 

643 

644 root: dict[TaskLabel, dict[str, dict[str | None, list[ProvenanceQuantumReport]]]] = {} 

645 

646 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

647 # when we inherit those docstrings in our public classes. 

648 if "sphinx" in sys.modules and not TYPE_CHECKING: 

649 

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

651 """See `pydantic.BaseModel.copy`.""" 

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

653 

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

655 """See `pydantic.BaseModel.model_dump`.""" 

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

657 

658 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

659 """See `pydantic.BaseModel.model_dump_json`.""" 

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

661 

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

663 """See `pydantic.BaseModel.model_copy`.""" 

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

665 

666 @classmethod 

667 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

668 """See `pydantic.BaseModel.model_construct`.""" 

669 return super().model_construct(*args, **kwargs) 

670 

671 @classmethod 

672 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

673 """See `pydantic.BaseModel.model_json_schema`.""" 

674 return super().model_json_schema(*args, **kwargs) 

675 

676 @classmethod 

677 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

678 """See `pydantic.BaseModel.model_validate`.""" 

679 return super().model_validate(*args, **kwargs) 

680 

681 @classmethod 

682 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

683 """See `pydantic.BaseModel.model_validate_json`.""" 

684 return super().model_validate_json(*args, **kwargs) 

685 

686 @classmethod 

687 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

688 """See `pydantic.BaseModel.model_validate_strings`.""" 

689 return super().model_validate_strings(*args, **kwargs) 

690 

691 

692class ProvenanceQuantumModel(pydantic.BaseModel): 

693 """Data model for the quanta in a provenance quantum graph file.""" 

694 

695 quantum_id: uuid.UUID 

696 """Unique identifier for the quantum.""" 

697 

698 task_label: TaskLabel 

699 """Name of the type of this dataset.""" 

700 

701 data_coordinate: DataCoordinateValues = pydantic.Field(default_factory=list) 

702 """The full values (required and implied) of this dataset's data ID.""" 

703 

704 inputs: dict[ConnectionName, list[uuid.UUID]] = pydantic.Field(default_factory=dict) 

705 """IDs of the datasets predicted to be consumed by this quantum, grouped by 

706 connection name. 

707 """ 

708 

709 outputs: dict[ConnectionName, list[uuid.UUID]] = pydantic.Field(default_factory=dict) 

710 """IDs of the datasets predicted to be produced by this quantum, grouped by 

711 connection name. 

712 """ 

713 

714 attempts: list[ProvenanceQuantumAttemptModel] = pydantic.Field(default_factory=list) 

715 """Provenance for all attempts to execute this quantum, ordered 

716 chronologically from first to last. 

717 

718 An entry is added merely if the quantum *should* have been attempted; an 

719 empty `list` is used only for quanta that were blocked by an upstream 

720 failure. 

721 """ 

722 

723 @property 

724 def node_id(self) -> uuid.UUID: 

725 """Alias for the quantum ID.""" 

726 return self.quantum_id 

727 

728 @classmethod 

729 def from_predicted(cls, predicted: PredictedQuantumDatasetsModel) -> ProvenanceQuantumModel: 

730 """Construct from a predicted quantum model. 

731 

732 Parameters 

733 ---------- 

734 predicted : `PredictedQuantumDatasetsModel` 

735 Information about the quantum from the predicted graph. 

736 

737 Returns 

738 ------- 

739 provenance : `ProvenanceQuantumModel` 

740 Provenance quantum model. 

741 """ 

742 inputs = { 

743 connection_name: [d.dataset_id for d in predicted_inputs] 

744 for connection_name, predicted_inputs in predicted.inputs.items() 

745 } 

746 outputs = { 

747 connection_name: [d.dataset_id for d in predicted_outputs] 

748 for connection_name, predicted_outputs in predicted.outputs.items() 

749 } 

750 return cls( 

751 quantum_id=predicted.quantum_id, 

752 task_label=predicted.task_label, 

753 data_coordinate=predicted.data_coordinate, 

754 inputs=inputs, 

755 outputs=outputs, 

756 ) 

757 

758 def _add_to_graph(self, graph: ProvenanceQuantumGraph) -> None: 

759 """Add this quantum and its edges to datasets to a provenance graph. 

760 

761 Parameters 

762 ---------- 

763 graph : `ProvenanceQuantumGraph` 

764 Graph to update in place. 

765 

766 Notes 

767 ----- 

768 This method adds: 

769 

770 - a ``bipartite_xgraph`` quantum node with full attributes; 

771 - a ``quantum_only_xgraph`` quantum node with full attributes; 

772 - ``bipartite_xgraph`` edges to adjacent datasets (which adds datasets 

773 nodes with no attributes), while populating those edge attributes; 

774 - ``quantum_only_xgraph`` edges to any adjacent quantum that has also 

775 already been loaded. 

776 """ 

777 task_node = graph.pipeline_graph.tasks[self.task_label] 

778 data_id = DataCoordinate.from_full_values(task_node.dimensions, tuple(self.data_coordinate)) 

779 last_attempt = ( 

780 self.attempts[-1] 

781 if self.attempts 

782 else ProvenanceQuantumAttemptModel(status=QuantumAttemptStatus.BLOCKED) 

783 ) 

784 graph._bipartite_xgraph.add_node( 

785 self.quantum_id, 

786 data_id=data_id, 

787 task_label=self.task_label, 

788 pipeline_node=task_node, 

789 status=last_attempt.status, 

790 caveats=last_attempt.caveats, 

791 exception=last_attempt.exception, 

792 resource_usage=last_attempt.resource_usage, 

793 attempts=self.attempts, 

794 ) 

795 graph._quanta_by_task_label[self.task_label][data_id] = self.quantum_id 

796 graph._quantum_only_xgraph.add_node(self.quantum_id, **graph._bipartite_xgraph.nodes[self.quantum_id]) 

797 for connection_name, dataset_ids in self.inputs.items(): 

798 read_edge = task_node.get_input_edge(connection_name) 

799 for dataset_id in dataset_ids: 

800 graph._bipartite_xgraph.add_edge(dataset_id, self.quantum_id, is_read=True) 

801 graph._bipartite_xgraph.edges[dataset_id, self.quantum_id].setdefault( 

802 "pipeline_edges", [] 

803 ).append(read_edge) 

804 for connection_name, dataset_ids in self.outputs.items(): 

805 write_edge = task_node.get_output_edge(connection_name) 

806 if connection_name == acc.METADATA_OUTPUT_CONNECTION_NAME: 

807 graph._bipartite_xgraph.add_node( 

808 dataset_ids[0], 

809 data_id=data_id, 

810 dataset_type_name=write_edge.dataset_type_name, 

811 pipeline_node=graph.pipeline_graph.dataset_types[write_edge.dataset_type_name], 

812 run=graph.header.output_run, 

813 produced=last_attempt.status.has_metadata, 

814 ) 

815 graph._datasets_by_type[write_edge.dataset_type_name][data_id] = dataset_ids[0] 

816 graph._bipartite_xgraph.nodes[self.quantum_id]["metadata_id"] = dataset_ids[0] 

817 graph._quantum_only_xgraph.nodes[self.quantum_id]["metadata_id"] = dataset_ids[0] 

818 if connection_name == acc.LOG_OUTPUT_CONNECTION_NAME: 

819 graph._bipartite_xgraph.add_node( 

820 dataset_ids[0], 

821 data_id=data_id, 

822 dataset_type_name=write_edge.dataset_type_name, 

823 pipeline_node=graph.pipeline_graph.dataset_types[write_edge.dataset_type_name], 

824 run=graph.header.output_run, 

825 produced=last_attempt.status.has_log, 

826 ) 

827 graph._datasets_by_type[write_edge.dataset_type_name][data_id] = dataset_ids[0] 

828 graph._bipartite_xgraph.nodes[self.quantum_id]["log_id"] = dataset_ids[0] 

829 graph._quantum_only_xgraph.nodes[self.quantum_id]["log_id"] = dataset_ids[0] 

830 for dataset_id in dataset_ids: 

831 graph._bipartite_xgraph.add_edge( 

832 self.quantum_id, 

833 dataset_id, 

834 is_read=False, 

835 # There can only be one pipeline edge for an output. 

836 pipeline_edges=[write_edge], 

837 ) 

838 for dataset_id in graph._bipartite_xgraph.predecessors(self.quantum_id): 

839 for upstream_quantum_id in graph._bipartite_xgraph.predecessors(dataset_id): 

840 graph._quantum_only_xgraph.add_edge(upstream_quantum_id, self.quantum_id) 

841 for dataset_id in graph._bipartite_xgraph.successors(self.quantum_id): 

842 for downstream_quantum_id in graph._bipartite_xgraph.successors(dataset_id): 

843 graph._quantum_only_xgraph.add_edge(self.quantum_id, downstream_quantum_id) 

844 

845 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

846 # when we inherit those docstrings in our public classes. 

847 if "sphinx" in sys.modules and not TYPE_CHECKING: 

848 

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

850 """See `pydantic.BaseModel.copy`.""" 

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

852 

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

854 """See `pydantic.BaseModel.model_dump`.""" 

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

856 

857 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

858 """See `pydantic.BaseModel.model_dump_json`.""" 

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

860 

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

862 """See `pydantic.BaseModel.model_copy`.""" 

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

864 

865 @classmethod 

866 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

867 """See `pydantic.BaseModel.model_construct`.""" 

868 return super().model_construct(*args, **kwargs) 

869 

870 @classmethod 

871 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

872 """See `pydantic.BaseModel.model_json_schema`.""" 

873 return super().model_json_schema(*args, **kwargs) 

874 

875 @classmethod 

876 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

877 """See `pydantic.BaseModel.model_validate`.""" 

878 return super().model_validate(*args, **kwargs) 

879 

880 @classmethod 

881 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

882 """See `pydantic.BaseModel.model_validate_json`.""" 

883 return super().model_validate_json(*args, **kwargs) 

884 

885 @classmethod 

886 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

887 """See `pydantic.BaseModel.model_validate_strings`.""" 

888 return super().model_validate_strings(*args, **kwargs) 

889 

890 

891class ProvenanceInitQuantumModel(pydantic.BaseModel): 

892 """Data model for the special "init" quanta in a provenance quantum graph 

893 file. 

894 """ 

895 

896 quantum_id: uuid.UUID 

897 """Unique identifier for the quantum.""" 

898 

899 task_label: TaskLabel 

900 """Name of the type of this dataset. 

901 

902 This is always a parent dataset type name, not a component. 

903 

904 Note that full dataset type definitions are stored in the pipeline graph. 

905 """ 

906 

907 inputs: dict[ConnectionName, uuid.UUID] = pydantic.Field(default_factory=dict) 

908 """IDs of the datasets predicted to be consumed by this quantum, grouped by 

909 connection name. 

910 """ 

911 

912 outputs: dict[ConnectionName, uuid.UUID] = pydantic.Field(default_factory=dict) 

913 """IDs of the datasets predicted to be produced by this quantum, grouped by 

914 connection name. 

915 """ 

916 

917 @classmethod 

918 def from_predicted(cls, predicted: PredictedQuantumDatasetsModel) -> ProvenanceInitQuantumModel: 

919 """Construct from a predicted quantum model. 

920 

921 Parameters 

922 ---------- 

923 predicted : `PredictedQuantumDatasetsModel` 

924 Information about the quantum from the predicted graph. 

925 

926 Returns 

927 ------- 

928 provenance : `ProvenanceInitQuantumModel` 

929 Provenance init quantum model. 

930 """ 

931 inputs = { 

932 connection_name: predicted_inputs[0].dataset_id 

933 for connection_name, predicted_inputs in predicted.inputs.items() 

934 } 

935 outputs = { 

936 connection_name: predicted_outputs[0].dataset_id 

937 for connection_name, predicted_outputs in predicted.outputs.items() 

938 } 

939 return cls( 

940 quantum_id=predicted.quantum_id, 

941 task_label=predicted.task_label, 

942 inputs=inputs, 

943 outputs=outputs, 

944 ) 

945 

946 def _add_to_graph(self, graph: ProvenanceQuantumGraph, empty_data_id: DataCoordinate) -> None: 

947 """Add this quantum and its edges to datasets to a provenance graph. 

948 

949 Parameters 

950 ---------- 

951 graph : `ProvenanceQuantumGraph` 

952 Graph to update in place. 

953 empty_data_id : `lsst.daf.butler.DataCoordinate` 

954 The empty data ID for the appropriate dimension universe. 

955 

956 Notes 

957 ----- 

958 This method adds: 

959 

960 - a ``bipartite_xgraph`` quantum node with full attributes; 

961 - ``bipartite_xgraph`` edges to adjacent datasets (which adds datasets 

962 nodes with no attributes), while populating those edge attributes; 

963 """ 

964 task_init_node = graph.pipeline_graph.tasks[self.task_label].init 

965 graph._bipartite_xgraph.add_node( 

966 self.quantum_id, data_id=empty_data_id, task_label=self.task_label, pipeline_node=task_init_node 

967 ) 

968 for connection_name, dataset_id in self.inputs.items(): 

969 read_edge = task_init_node.get_input_edge(connection_name) 

970 graph._bipartite_xgraph.add_edge(dataset_id, self.quantum_id, is_read=True) 

971 graph._bipartite_xgraph.edges[dataset_id, self.quantum_id].setdefault( 

972 "pipeline_edges", [] 

973 ).append(read_edge) 

974 for connection_name, dataset_id in self.outputs.items(): 

975 write_edge = task_init_node.get_output_edge(connection_name) 

976 graph._bipartite_xgraph.add_node( 

977 dataset_id, 

978 data_id=empty_data_id, 

979 dataset_type_name=write_edge.dataset_type_name, 

980 pipeline_node=graph.pipeline_graph.dataset_types[write_edge.dataset_type_name], 

981 run=graph.header.output_run, 

982 produced=True, 

983 ) 

984 graph._datasets_by_type[write_edge.dataset_type_name][empty_data_id] = dataset_id 

985 graph._bipartite_xgraph.add_edge( 

986 self.quantum_id, 

987 dataset_id, 

988 is_read=False, 

989 # There can only be one pipeline edge for an output. 

990 pipeline_edges=[write_edge], 

991 ) 

992 if write_edge.connection_name == acc.CONFIG_INIT_OUTPUT_CONNECTION_NAME: 

993 graph._bipartite_xgraph.nodes[self.quantum_id]["config_id"] = dataset_id 

994 graph._init_quanta[self.task_label] = self.quantum_id 

995 

996 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

997 # when we inherit those docstrings in our public classes. 

998 if "sphinx" in sys.modules and not TYPE_CHECKING: 

999 

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

1001 """See `pydantic.BaseModel.copy`.""" 

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

1003 

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

1005 """See `pydantic.BaseModel.model_dump`.""" 

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

1007 

1008 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

1009 """See `pydantic.BaseModel.model_dump_json`.""" 

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

1011 

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

1013 """See `pydantic.BaseModel.model_copy`.""" 

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

1015 

1016 @classmethod 

1017 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

1018 """See `pydantic.BaseModel.model_construct`.""" 

1019 return super().model_construct(*args, **kwargs) 

1020 

1021 @classmethod 

1022 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

1023 """See `pydantic.BaseModel.model_json_schema`.""" 

1024 return super().model_json_schema(*args, **kwargs) 

1025 

1026 @classmethod 

1027 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

1028 """See `pydantic.BaseModel.model_validate`.""" 

1029 return super().model_validate(*args, **kwargs) 

1030 

1031 @classmethod 

1032 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

1033 """See `pydantic.BaseModel.model_validate_json`.""" 

1034 return super().model_validate_json(*args, **kwargs) 

1035 

1036 @classmethod 

1037 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

1038 """See `pydantic.BaseModel.model_validate_strings`.""" 

1039 return super().model_validate_strings(*args, **kwargs) 

1040 

1041 

1042class ProvenanceInitQuantaModel(pydantic.RootModel): 

1043 """Data model for the init quanta in a provenance graph.""" 

1044 

1045 root: list[ProvenanceInitQuantumModel] = pydantic.Field(default_factory=list) 

1046 """List of special "init" quanta, one for each task.""" 

1047 

1048 def _add_to_graph(self, graph: ProvenanceQuantumGraph) -> None: 

1049 """Add this quantum and its edges to datasets to a provenance graph. 

1050 

1051 Parameters 

1052 ---------- 

1053 graph : `ProvenanceQuantumGraph` 

1054 Graph to update in place. 

1055 """ 

1056 empty_data_id = DataCoordinate.make_empty(graph.pipeline_graph.universe) 

1057 for init_quantum in self.root: 

1058 init_quantum._add_to_graph(graph, empty_data_id=empty_data_id) 

1059 

1060 # Work around the fact that Sphinx chokes on Pydantic docstring formatting, 

1061 # when we inherit those docstrings in our public classes. 

1062 if "sphinx" in sys.modules and not TYPE_CHECKING: 

1063 

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

1065 """See `pydantic.BaseModel.copy`.""" 

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

1067 

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

1069 """See `pydantic.BaseModel.model_dump`.""" 

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

1071 

1072 def model_dump_json(self, *args: Any, **kwargs: Any) -> Any: 

1073 """See `pydantic.BaseModel.model_dump_json`.""" 

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

1075 

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

1077 """See `pydantic.BaseModel.model_copy`.""" 

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

1079 

1080 @classmethod 

1081 def model_construct(cls, *args: Any, **kwargs: Any) -> Any: # type: ignore[misc, override] 

1082 """See `pydantic.BaseModel.model_construct`.""" 

1083 return super().model_construct(*args, **kwargs) 

1084 

1085 @classmethod 

1086 def model_json_schema(cls, *args: Any, **kwargs: Any) -> Any: 

1087 """See `pydantic.BaseModel.model_json_schema`.""" 

1088 return super().model_json_schema(*args, **kwargs) 

1089 

1090 @classmethod 

1091 def model_validate(cls, *args: Any, **kwargs: Any) -> Any: 

1092 """See `pydantic.BaseModel.model_validate`.""" 

1093 return super().model_validate(*args, **kwargs) 

1094 

1095 @classmethod 

1096 def model_validate_json(cls, *args: Any, **kwargs: Any) -> Any: 

1097 """See `pydantic.BaseModel.model_validate_json`.""" 

1098 return super().model_validate_json(*args, **kwargs) 

1099 

1100 @classmethod 

1101 def model_validate_strings(cls, *args: Any, **kwargs: Any) -> Any: 

1102 """See `pydantic.BaseModel.model_validate_strings`.""" 

1103 return super().model_validate_strings(*args, **kwargs) 

1104 

1105 

1106class ProvenanceQuantumGraph(BaseQuantumGraph): 

1107 """A quantum graph that represents processing that has already been 

1108 executed. 

1109 

1110 Parameters 

1111 ---------- 

1112 header : `HeaderModel` 

1113 General metadata shared with other quantum graph types. 

1114 pipeline_graph : `.pipeline_graph.PipelineGraph` 

1115 Graph of tasks and dataset types. May contain a superset of the tasks 

1116 and dataset types that actually have quanta and datasets in the quantum 

1117 graph. 

1118 

1119 Notes 

1120 ----- 

1121 A provenance quantum graph is generally obtained via the 

1122 `ProvenanceQuantumGraphReader.graph` attribute, which is updated in-place 

1123 as information is read from disk. 

1124 """ 

1125 

1126 def __init__(self, header: HeaderModel, pipeline_graph: PipelineGraph) -> None: 

1127 super().__init__(header, pipeline_graph) 

1128 self._init_quanta: dict[TaskLabel, uuid.UUID] = {} 

1129 self._quantum_only_xgraph = networkx.DiGraph() 

1130 self._bipartite_xgraph = networkx.DiGraph() 

1131 self._quanta_by_task_label: dict[str, dict[DataCoordinate, uuid.UUID]] = { 

1132 task_label: {} for task_label in self.pipeline_graph.tasks.keys() 

1133 } 

1134 self._datasets_by_type: dict[str, dict[DataCoordinate, uuid.UUID]] = { 

1135 dataset_type_name: {} for dataset_type_name in self.pipeline_graph.dataset_types.keys() 

1136 } 

1137 

1138 @classmethod 

1139 @contextmanager 

1140 def from_args( 

1141 cls, 

1142 repo_or_filename: str, 

1143 /, 

1144 collection: str | None = None, 

1145 *, 

1146 quanta: Iterable[uuid.UUID] | None = None, 

1147 datasets: Iterable[uuid.UUID] | None = None, 

1148 writeable: bool = False, 

1149 ) -> Iterator[tuple[ProvenanceQuantumGraph, Butler | None]]: 

1150 """Construct a `ProvenanceQuantumGraph` fron CLI-friendly arguments for 

1151 a file or butler-ingested graph dataset. 

1152 

1153 Parameters 

1154 ---------- 

1155 repo_or_filename : `str` 

1156 Either a provenance quantum graph filename or a butler repository 

1157 path or alias. 

1158 collection : `str`, optional 

1159 Collection to search; presence indicates that the first argument 

1160 is a butler repository, not a filename. 

1161 quanta : `~collections.abc.Iterable` [ `str` ] or `None`, optional 

1162 IDs of the quanta to load, or `None` to load all. 

1163 datasets : `~collections.abc.Iterable` [ `str` ], optional 

1164 IDs of the datasets to load, or `None` to load all. 

1165 writeable : `bool`, optional 

1166 Whether the butler should be constructed with write support. 

1167 

1168 Returns 

1169 ------- 

1170 context : `contextlib.AbstractContextManager` 

1171 A context manager that yields a tuple of 

1172 

1173 - the `ProvenanceQuantumGraph` 

1174 - the `Butler` constructed (or `None`) 

1175 

1176 when entered. 

1177 """ 

1178 exit_stack = ExitStack() 

1179 if collection is not None: 

1180 try: 

1181 butler = exit_stack.enter_context( 

1182 Butler.from_config(repo_or_filename, collections=[collection], writeable=writeable) 

1183 ) 

1184 except Exception as err: 

1185 err.add_note( 

1186 f"Expected {repo_or_filename!r} to be a butler repository path or alias because a " 

1187 f"collection ({collection}) was provided." 

1188 ) 

1189 raise 

1190 with exit_stack: 

1191 graph = butler.get( 

1192 acc.PROVENANCE_DATASET_TYPE_NAME, parameters={"quanta": quanta, "datasets": datasets} 

1193 ) 

1194 yield graph, butler 

1195 else: 

1196 try: 

1197 reader = exit_stack.enter_context(ProvenanceQuantumGraphReader.open(repo_or_filename)) 

1198 except Exception as err: 

1199 err.add_note( 

1200 f"Expected a {repo_or_filename} to be a provenance quantum graph filename " 

1201 f"because no collection was provided." 

1202 ) 

1203 raise 

1204 with exit_stack: 

1205 if quanta is None: 1205 ↛ 1207line 1205 didn't jump to line 1207 because the condition on line 1205 was always true

1206 reader.read_quanta() 

1207 elif not quanta: 

1208 reader.read_quanta(quanta) 

1209 if datasets is None: 1209 ↛ 1210line 1209 didn't jump to line 1210 because the condition on line 1209 was never true

1210 reader.read_datasets() 

1211 elif not datasets: 1211 ↛ 1213line 1211 didn't jump to line 1213 because the condition on line 1211 was always true

1212 reader.read_datasets(datasets) 

1213 yield reader.graph, None 

1214 

1215 @property 

1216 def init_quanta(self) -> Mapping[TaskLabel, uuid.UUID]: 

1217 """A mapping from task label to the ID of the special init quantum for 

1218 that task. 

1219 

1220 This is populated by the ``init_quanta`` component. Additional 

1221 information about each init quantum can be found by using the ID to 

1222 look up node attributes in the `bipartite_xgraph`, i.e.:: 

1223 

1224 info: ProvenanceInitQuantumInfo = qg.bipartite_xgraph.nodes[id] 

1225 """ 

1226 return self._init_quanta 

1227 

1228 @property 

1229 def quanta_by_task(self) -> Mapping[TaskLabel, Mapping[DataCoordinate, uuid.UUID]]: 

1230 """A nested mapping of all quanta, keyed first by task name and then by 

1231 data ID. 

1232 

1233 Notes 

1234 ----- 

1235 This is populated one quantum at a time as they are read. All tasks in 

1236 the pipeline graph are included, even if none of their quanta were 

1237 loaded (i.e. nested mappings may be empty). 

1238 

1239 The returned object may be an internal dictionary; as the type 

1240 annotation indicates, it should not be modified in place. 

1241 """ 

1242 return self._quanta_by_task_label 

1243 

1244 @property 

1245 def datasets_by_type(self) -> Mapping[DatasetTypeName, Mapping[DataCoordinate, uuid.UUID]]: 

1246 """A nested mapping of all datasets, keyed first by dataset type name 

1247 and then by data ID. 

1248 

1249 Notes 

1250 ----- 

1251 This is populated one dataset at a time as they are read. All dataset 

1252 types in the pipeline graph are included, even if none of their 

1253 datasets were loaded (i.e. nested mappings may be empty). 

1254 

1255 Reading a quantum also populates its log and metadata datasets. 

1256 

1257 The returned object may be an internal dictionary; as the type 

1258 annotation indicates, it should not be modified in place. 

1259 """ 

1260 return self._datasets_by_type 

1261 

1262 @property 

1263 def quantum_only_xgraph(self) -> networkx.DiGraph: 

1264 """A directed acyclic graph with quanta as nodes (and datasets elided). 

1265 

1266 Notes 

1267 ----- 

1268 Node keys are quantum UUIDs, and are populated one quantum at a time as 

1269 they are loaded. Loading quanta (via 

1270 `ProvenanceQuantumGraphReader.read_quanta`) will add the loaded nodes 

1271 with full attributes and add edges to adjacent nodes with no 

1272 attributes. Loading datasets (via 

1273 `ProvenanceQuantumGraphReader.read_datasets`) will also add edges and 

1274 nodes with no attributes. 

1275 

1276 Node attributes are described by the `ProvenanceQuantumInfo` types. 

1277 

1278 This graph does not include special "init" quanta. 

1279 

1280 The returned object is a read-only view of an internal one. 

1281 """ 

1282 return self._quantum_only_xgraph.copy(as_view=True) 

1283 

1284 @property 

1285 def bipartite_xgraph(self) -> networkx.DiGraph: 

1286 """A directed acyclic graph with quantum and dataset nodes. 

1287 

1288 Notes 

1289 ----- 

1290 Node keys are quantum or dataset UUIDs, and are populated one quantum 

1291 or dataset at a time as they are loaded. Loading quanta (via 

1292 `ProvenanceQuantumGraphReader.read_quanta`) or datasets (via 

1293 `ProvenanceQuantumGraphReader.read_datasets`) will load those nodes 

1294 with full attributes and edges to adjacent nodes with no attributes. 

1295 Loading quanta is necessary to populate edge attributes. 

1296 Reading a quantum also populates its log and metadata datasets. 

1297 

1298 Node attributes are described by the 

1299 `ProvenanceQuantumInfo`, `ProvenanceInitQuantumInfo`, and 

1300 `ProvenanceDatasetInfo` types. 

1301 

1302 This graph includes init-input and init-output datasets, but it does 

1303 *not* reflect the dependency between each task's special "init" quantum 

1304 and its runtime quanta (as this would require edges between quanta, and 

1305 that would break the "bipartite" property). 

1306 

1307 The returned object is a read-only view of an internal one. 

1308 """ 

1309 return self._bipartite_xgraph.copy(as_view=True) 

1310 

1311 def make_quantum_table( 

1312 self, *, drop_unused_columns: bool = True, expand_caveats: bool = False, as_table: bool = True 

1313 ) -> astropy.table.Table | list[dict]: 

1314 """Construct an `astropy.table.Table` with a tabular summary of the 

1315 quanta. 

1316 

1317 Parameters 

1318 ---------- 

1319 drop_unused_columns : `bool`, optional 

1320 Whether to drop columns for rare states that did not actually 

1321 occur in this run. 

1322 

1323 expand_caveats : `bool`, optional 

1324 Whether to display a comma-separated list of task caveats instead 

1325 of a collapsed `multiple` marker. 

1326 

1327 as_table : `bool`, optional 

1328 Whether to return an `astropy.table.Table` or a list of rows. 

1329 

1330 Returns 

1331 ------- 

1332 table : `astropy.table.Table` | `list` [`dict`] 

1333 A table view of the quantum information. This only includes 

1334 counts of status categories and caveats, not any per-data-ID 

1335 detail. If `as_table` is False, returns a raw list of rows. 

1336 

1337 Notes 

1338 ----- 

1339 Success caveats in the table are represented by their 

1340 `~QuantumSuccessCaveats.concise` form, so when pretty-printing this 

1341 table for users, the `~QuantumSuccessCaveats.legend` should generally 

1342 be printed as well. 

1343 """ 

1344 rows: list[dict] = [] 

1345 for task_label, quanta_for_task in self.quanta_by_task.items(): 

1346 if not self.header.n_task_quanta[task_label]: 1346 ↛ 1347line 1346 didn't jump to line 1347 because the condition on line 1346 was never true

1347 continue 

1348 status_counts: Counter[QuantumAttemptStatus] = Counter( 

1349 self._quantum_only_xgraph.nodes[q]["status"] for q in quanta_for_task.values() 

1350 ) 

1351 caveat_counts: Counter[QuantumSuccessCaveats | None] = Counter( 

1352 self._quantum_only_xgraph.nodes[q]["caveats"] for q in quanta_for_task.values() 

1353 ) 

1354 caveat_counts.pop(QuantumSuccessCaveats.NO_CAVEATS, None) 

1355 caveat_counts.pop(None, None) 

1356 

1357 if len(caveat_counts) > 1 and expand_caveats: 1357 ↛ 1358line 1357 didn't jump to line 1358 because the condition on line 1357 was never true

1358 caveats = ",".join( 

1359 f"{code.concise()}({count})" for code, count in caveat_counts.items() if code is not None 

1360 ) 

1361 elif len(caveat_counts) > 1: 1361 ↛ 1362line 1361 didn't jump to line 1362 because the condition on line 1361 was never true

1362 caveats = "(multiple)" 

1363 elif len(caveat_counts) == 1: 

1364 ((code, count),) = caveat_counts.items() 

1365 if TYPE_CHECKING: 

1366 assert code is not None 

1367 caveats = f"{code.concise()}({count})" 

1368 else: 

1369 caveats = "" 

1370 row: dict[str, Any] = { 

1371 "Task": task_label, 

1372 "Caveats": caveats, 

1373 } 

1374 for status in QuantumAttemptStatus: 

1375 row[status.title] = status_counts.get(status, 0) 

1376 row.update( 

1377 { 

1378 "TOTAL": len(quanta_for_task), 

1379 "EXPECTED": self.header.n_task_quanta[task_label], 

1380 } 

1381 ) 

1382 rows.append(row) 

1383 if not as_table: 

1384 return rows 

1385 table = astropy.table.Table(rows) 

1386 if drop_unused_columns: 1386 ↛ 1390line 1386 didn't jump to line 1390 because the condition on line 1386 was always true

1387 for status in QuantumAttemptStatus: 

1388 if status.is_rare and not table[status.title].any(): 

1389 del table[status.title] 

1390 return table 

1391 

1392 def make_exception_table(self, *, as_table: bool = True) -> astropy.table.Table | list[dict]: 

1393 """Construct an `astropy.table.Table` with counts for each exception 

1394 type raised by each task. 

1395 

1396 Parameters 

1397 ---------- 

1398 as_table : `bool`, optional 

1399 Whether to return an `astropy.table.Table` or a list of rows. 

1400 

1401 Returns 

1402 ------- 

1403 table : `astropy.table.Table` | `list` [`dict`] 

1404 A table with columns for task label, exception type, and counts, 

1405 or the raw table rows as a list of dicts. 

1406 """ 

1407 rows: list[dict] = [] 

1408 for task_label, quanta_for_task in self.quanta_by_task.items(): 

1409 success_counts = Counter[str]() 

1410 failed_counts = Counter[str]() 

1411 for quantum_id in quanta_for_task.values(): 

1412 quantum_info: ProvenanceQuantumInfo = self._quantum_only_xgraph.nodes[quantum_id] 

1413 exc_info = quantum_info["exception"] 

1414 if exc_info is not None: 

1415 if quantum_info["status"] is QuantumAttemptStatus.SUCCESSFUL: 

1416 success_counts[exc_info.type_name] += 1 

1417 else: 

1418 failed_counts[exc_info.type_name] += 1 

1419 for type_name in sorted(success_counts.keys() | failed_counts.keys()): 

1420 rows.append( 

1421 { 

1422 "Task": task_label, 

1423 "Exception": type_name, 

1424 "Successes": success_counts.get(type_name, 0), 

1425 "Failures": failed_counts.get(type_name, 0), 

1426 } 

1427 ) 

1428 if as_table: 

1429 return astropy.table.Table(rows) 

1430 else: 

1431 return rows 

1432 

1433 def make_task_resource_usage_table( 

1434 self, task_label: TaskLabel, include_data_ids: bool = False 

1435 ) -> astropy.table.Table: 

1436 """Make a table of resource usage for a single task. 

1437 

1438 Parameters 

1439 ---------- 

1440 task_label : `str` 

1441 Label of the task to extract resource usage for. 

1442 include_data_ids : `bool`, optional 

1443 Whether to also include data ID columns. 

1444 

1445 Returns 

1446 ------- 

1447 table : `astropy.table.Table` 

1448 A table with columns for quantum ID and all fields in 

1449 `QuantumResourceUsage`. 

1450 """ 

1451 quanta_for_task = self.quanta_by_task[task_label] 

1452 dtype_terms: list[tuple[str, np.dtype]] = [("quantum_id", np.dtype((np.void, 16)))] 

1453 if include_data_ids: 1453 ↛ 1458line 1453 didn't jump to line 1458 because the condition on line 1453 was always true

1454 dimensions = self.pipeline_graph.tasks[task_label].dimensions 

1455 for dimension_name in dimensions.data_coordinate_keys: 

1456 dtype = np.dtype(self.pipeline_graph.universe.dimensions[dimension_name].primary_key.pytype) 

1457 dtype_terms.append((dimension_name, dtype)) 

1458 fields = QuantumResourceUsage.get_numpy_fields() 

1459 dtype_terms.extend(fields.items()) 

1460 row_dtype = np.dtype(dtype_terms) 

1461 rows: list[object] = [] 

1462 for data_id, quantum_id in quanta_for_task.items(): 

1463 info: ProvenanceQuantumInfo = self._quantum_only_xgraph.nodes[quantum_id] 

1464 if (resource_usage := info["resource_usage"]) is not None: 1464 ↛ 1462line 1464 didn't jump to line 1462 because the condition on line 1464 was always true

1465 row: tuple[object, ...] = (quantum_id.bytes,) 

1466 if include_data_ids: 1466 ↛ 1468line 1466 didn't jump to line 1468 because the condition on line 1466 was always true

1467 row += data_id.full_values 

1468 row += resource_usage.get_numpy_row() 

1469 rows.append(row) 

1470 array = np.array(rows, dtype=row_dtype) 

1471 return astropy.table.Table(array, units=QuantumResourceUsage.get_units()) 

1472 

1473 def make_status_report( 

1474 self, 

1475 states: Iterable[QuantumAttemptStatus] = ( 

1476 QuantumAttemptStatus.FAILED, 

1477 QuantumAttemptStatus.ABORTED, 

1478 QuantumAttemptStatus.ABORTED_SUCCESS, 

1479 ), 

1480 *, 

1481 also: QuantumAttemptStatus | Iterable[QuantumAttemptStatus] = (), 

1482 with_caveats: QuantumSuccessCaveats | None = QuantumSuccessCaveats.PARTIAL_OUTPUTS_ERROR, 

1483 data_id_table_dir: ResourcePathExpression | None = None, 

1484 ) -> ProvenanceReport: 

1485 """Make a JSON- or YAML-friendly report of all quanta with the given 

1486 states. 

1487 

1488 Parameters 

1489 ---------- 

1490 states : `~collections.abc.Iterable` [`..QuantumAttemptStatus`] or \ 

1491 `..QuantumAttemptStatus`, optional 

1492 A quantum is included if it has any of these states. Defaults to 

1493 states that clearly represent problems. 

1494 also : `~collections.abc.Iterable` [`..QuantumAttemptStatus`] or \ 

1495 `..QuantumAttemptStatus`, optional 

1496 Additional states to consider; unioned with ``states``. This is 

1497 provided so users can easily request additional states while also 

1498 getting the defaults. 

1499 with_caveats : `..QuantumSuccessCaveats` or `None`, optional 

1500 If `..QuantumAttemptStatus.SUCCESSFUL` is in ``states``, only 

1501 include quanta with these caveat flags. May be set to `None` 

1502 to report on all successful quanta. 

1503 data_id_table_dir : convertible to `~lsst.resources.ResourcePath`, \ 

1504 optional 

1505 If provided, a directory to write data ID tables (in ECSV format) 

1506 with all of the data IDs with the given states, for use with the 

1507 ``--data-id-tables`` argument to the quantum graph builder. 

1508 Subdirectories for each task and status will created within this 

1509 directory, with one file for each exception type (or ``UNKNOWN`` 

1510 when there is no exception). 

1511 

1512 Returns 

1513 ------- 

1514 report : `ProvenanceModel` 

1515 A Pydantic model that groups quanta by task label and exception 

1516 type. 

1517 """ 

1518 states = set(ensure_iterable(states)) 

1519 states.update(ensure_iterable(also)) 

1520 result = ProvenanceReport(root={}) 

1521 if data_id_table_dir is not None: 1521 ↛ 1523line 1521 didn't jump to line 1523 because the condition on line 1521 was always true

1522 data_id_table_dir = ResourcePath(data_id_table_dir) 

1523 for task_label, quanta_for_task in self.quanta_by_task.items(): 

1524 reports_for_task: dict[str, dict[str | None, list[ProvenanceQuantumReport]]] = {} 

1525 table_rows_for_task: dict[str, dict[str | None, list[tuple[int | str, ...]]]] = {} 

1526 for quantum_id in quanta_for_task.values(): 

1527 quantum_info: ProvenanceQuantumInfo = self._quantum_only_xgraph.nodes[quantum_id] 

1528 quantum_status = quantum_info["status"] 

1529 if quantum_status not in states: 

1530 continue 

1531 if ( 

1532 quantum_status is QuantumAttemptStatus.SUCCESSFUL 

1533 and with_caveats is not None 

1534 and (quantum_info["caveats"] is None or not (quantum_info["caveats"] & with_caveats)) 

1535 ): 

1536 continue 

1537 key1 = quantum_status.name 

1538 exc_info = quantum_info["exception"] 

1539 key2 = exc_info.type_name if exc_info is not None else None 

1540 reports_for_task.setdefault(key1, {}).setdefault(key2, []).append( 

1541 ProvenanceQuantumReport.from_info(quantum_id, quantum_info) 

1542 ) 

1543 if data_id_table_dir: 1543 ↛ 1526line 1543 didn't jump to line 1526 because the condition on line 1543 was always true

1544 table_rows_for_task.setdefault(key1, {}).setdefault(key2, []).append( 

1545 quantum_info["data_id"].required_values 

1546 ) 

1547 if reports_for_task: 

1548 result.root[task_label] = reports_for_task 

1549 if table_rows_for_task: 

1550 assert data_id_table_dir is not None, "table_rows_for_task should be empty" 

1551 for status_name, table_rows_for_status in table_rows_for_task.items(): 

1552 dir_for_task_and_status = data_id_table_dir.join(task_label, forceDirectory=True).join( 

1553 status_name, forceDirectory=True 

1554 ) 

1555 if dir_for_task_and_status.isLocal: 1555 ↛ 1557line 1555 didn't jump to line 1557 because the condition on line 1555 was always true

1556 dir_for_task_and_status.mkdir() 

1557 for exc_name, data_id_rows in table_rows_for_status.items(): 

1558 table = astropy.table.Table( 

1559 rows=data_id_rows, 

1560 names=list(self.pipeline_graph.tasks[task_label].dimensions.required), 

1561 ) 

1562 filename = f"{exc_name}.ecsv" if exc_name is not None else "UNKNOWN.ecsv" 

1563 with dir_for_task_and_status.join(filename).open("w") as stream: 

1564 table.write(stream, format="ecsv") 

1565 return result 

1566 

1567 def make_many_reports( 

1568 self, 

1569 states: Iterable[QuantumAttemptStatus] = ( 

1570 QuantumAttemptStatus.FAILED, 

1571 QuantumAttemptStatus.ABORTED, 

1572 QuantumAttemptStatus.ABORTED_SUCCESS, 

1573 ), 

1574 *, 

1575 status_report_file: ResourcePathExpression | None = None, 

1576 print_quantum_table: bool = False, 

1577 print_exception_table: bool = False, 

1578 also: QuantumAttemptStatus | Iterable[QuantumAttemptStatus] = (), 

1579 with_caveats: QuantumSuccessCaveats | None = None, 

1580 data_id_table_dir: ResourcePathExpression | None = None, 

1581 output_format: Literal["json", "table"] = "table", 

1582 print_legend: bool = True, 

1583 expand_caveats: bool = False, 

1584 ) -> None: 

1585 """Write multiple reports. 

1586 

1587 Parameters 

1588 ---------- 

1589 states : `~collections.abc.Iterable` [`..QuantumAttemptStatus`] or \ 

1590 `..QuantumAttemptStatus`, optional 

1591 A quantum is included in the status report and data ID tables if it 

1592 has any of these states. Defaults to states that clearly represent 

1593 problems. 

1594 status_report_file : convertible to `~lsst.resources.ResourcePath`, 

1595 optional 

1596 Filename for the JSON status report (see `make_status_report`). 

1597 print_quantum_table : `bool`, optional 

1598 If `True`, print a quantum summary table (counts only) to STDOUT. 

1599 print_exception_table : `bool`, optional 

1600 If `True`, print an exception-type summary table (counts only) to 

1601 STDOUT. 

1602 also : `~collections.abc.Iterable` [`..QuantumAttemptStatus`] or \ 

1603 `..QuantumAttemptStatus`, optional 

1604 Additional states to consider in the status report and data ID 

1605 tables; unioned with ``states``. This is provided so users can 

1606 easily request additional states while also getting the defaults. 

1607 with_caveats : `..QuantumSuccessCaveats` or `None`, optional 

1608 Only include quanta with these caveat flags in the status report 

1609 and data ID tables. May be set to `None` to report on all 

1610 successful quanta (an empty sequence reports on only quanta with no 

1611 caveats). If provided, `QuantumAttemptStatus.SUCCESSFUL` is 

1612 automatically included in ``states``. 

1613 data_id_table_dir : convertible to `~lsst.resources.ResourcePath`, \ 

1614 optional 

1615 If provided, a directory to write data ID tables (in ECSV format) 

1616 with all of the data IDs with the given states, for use with the 

1617 ``--data-id-tables`` argument to the quantum graph builder. 

1618 Subdirectories for each task and status will created within this 

1619 directory, with one file for each exception type (or ``UNKNOWN`` 

1620 when there is no exception). 

1621 output_format : `str`, optional 

1622 Sets the output format for the data printed to stdout. Defaults to 

1623 `table`, which presents one or more text tables; `json` produces 

1624 machine-readable JSON formatted data. 

1625 print_legend : `bool`, optional 

1626 When `True`, prints a legend section describing the caveats 

1627 notation after displaying the output tables. 

1628 expand_caveats : `bool`, optional 

1629 If `False` (default), multiple caveats will be reduced to a single 

1630 value; if `True`, then all available caveats will be displayed. 

1631 This option is set to `True` when the output format is JSON. 

1632 """ 

1633 as_table = output_format == "table" 

1634 

1635 if output_format == "json": 

1636 json_obj: dict[str, Any] = {"tasks": None, "exceptions": None} 

1637 expand_caveats = True 

1638 

1639 if TYPE_CHECKING: 

1640 assert json_obj 

1641 

1642 if status_report_file is not None or data_id_table_dir is not None: 

1643 status_report = self.make_status_report( 

1644 states, also=also, with_caveats=with_caveats, data_id_table_dir=data_id_table_dir 

1645 ) 

1646 if status_report_file is not None: 1646 ↛ 1652line 1646 didn't jump to line 1652 because the condition on line 1646 was always true

1647 status_report_file = ResourcePath(status_report_file) 

1648 if status_report_file.isLocal: 1648 ↛ 1650line 1648 didn't jump to line 1650 because the condition on line 1648 was always true

1649 status_report_file.dirname().mkdir() 

1650 with ResourcePath(status_report_file).open("w") as stream: 

1651 stream.write(status_report.model_dump_json(indent=2)) 

1652 if print_quantum_table: 1652 ↛ 1675line 1652 didn't jump to line 1675 because the condition on line 1652 was always true

1653 quantum_table = self.make_quantum_table(expand_caveats=expand_caveats, as_table=as_table) 

1654 match (len(quantum_table) > 0, quantum_table): 

1655 case (True, astropy.table.Table()): 

1656 quantum_table.pprint_all() 

1657 print("") 

1658 if print_legend: 1658 ↛ 1662line 1658 didn't jump to line 1662 because the condition on line 1658 was always true

1659 print("Caveats\n-------") 

1660 for k, v in QuantumSuccessCaveats.legend().items(): 

1661 print(f"{k}: {v}") 

1662 print("") 

1663 case (True, list()): 1663 ↛ 1673line 1663 didn't jump to line 1673 because the pattern on line 1663 always matched

1664 json_obj["tasks"] = [ 

1665 { 

1666 **row, 

1667 "Caveats": [ 

1668 QuantumSuccessCaveats.expanded_dict(w) for w in row["Caveats"].split(",") 

1669 ], 

1670 } 

1671 for row in quantum_table 

1672 ] 

1673 case _: 

1674 pass 

1675 if print_exception_table: 1675 ↛ 1686line 1675 didn't jump to line 1686 because the condition on line 1675 was always true

1676 exception_table = self.make_exception_table(as_table=as_table) 

1677 match (len(exception_table) > 0, exception_table): 

1678 case (True, astropy.table.Table()): 

1679 exception_table.pprint_all() 

1680 print("") 

1681 case (True, list()): 1681 ↛ 1683line 1681 didn't jump to line 1683 because the pattern on line 1681 always matched

1682 json_obj["exceptions"] = exception_table 

1683 case _: 

1684 pass 

1685 

1686 if output_format == "json": 

1687 json_obj["legend"] = QuantumSuccessCaveats.legend() 

1688 json.dump(json_obj, sys.stdout) 

1689 

1690 

1691@dataclasses.dataclass 

1692class ProvenanceQuantumGraphReader(BaseQuantumGraphReader): 

1693 """A helper class for reading provenance quantum graphs. 

1694 

1695 Notes 

1696 ----- 

1697 The `open` context manager should be used to construct new instances. 

1698 Instances cannot be used after the context manager exits, except to access 

1699 the `graph` attribute`. 

1700 

1701 The various ``read_*`` methods in this class update the `graph` attribute 

1702 in place. 

1703 """ 

1704 

1705 graph: ProvenanceQuantumGraph = dataclasses.field(init=False) 

1706 """Loaded provenance graph, populated in place as components are read.""" 

1707 

1708 @classmethod 

1709 @contextmanager 

1710 def open( 

1711 cls, 

1712 uri: ResourcePathExpression, 

1713 *, 

1714 page_size: int | None = None, 

1715 import_mode: TaskImportMode = TaskImportMode.DO_NOT_IMPORT, 

1716 ) -> Iterator[ProvenanceQuantumGraphReader]: 

1717 """Construct a reader from a URI. 

1718 

1719 Parameters 

1720 ---------- 

1721 uri : convertible to `lsst.resources.ResourcePath` 

1722 URI to open. Should have a ``.qg`` extension. 

1723 page_size : `int`, optional 

1724 Approximate number of bytes to read at once from address files and 

1725 multi-block files. Note that this does not set a page size for 

1726 *all* reads, but it does affect the smallest, most numerous reads. 

1727 Can also be set via the ``LSST_QG_PAGE_SIZE`` environment variable. 

1728 import_mode : `.pipeline_graph.TaskImportMode`, optional 

1729 How to handle importing the task classes referenced in the pipeline 

1730 graph. 

1731 

1732 Returns 

1733 ------- 

1734 reader : `contextlib.AbstractContextManager` [ \ 

1735 `ProvenanceQuantumGraphReader` ] 

1736 A context manager that returns the reader when entered. 

1737 """ 

1738 with cls._open( 

1739 uri, 

1740 graph_type="provenance", 

1741 address_filename="nodes", 

1742 page_size=page_size, 

1743 import_mode=import_mode, 

1744 n_addresses=4, 

1745 ) as self: 

1746 yield self 

1747 

1748 def __post_init__(self) -> None: 

1749 self.graph = ProvenanceQuantumGraph(self.header, self.pipeline_graph) 

1750 

1751 def read_init_quanta(self) -> None: 

1752 """Read the thin graph, with all edge information and categorization of 

1753 quanta by task label. 

1754 """ 

1755 init_quanta = self._read_single_block("init_quanta", ProvenanceInitQuantaModel) 

1756 for init_quantum in init_quanta.root: 

1757 self.graph._init_quanta[init_quantum.task_label] = init_quantum.quantum_id 

1758 init_quanta._add_to_graph(self.graph) 

1759 

1760 def read_full_graph(self) -> None: 

1761 """Read all bipartite edges and all quantum and dataset node 

1762 attributes, fully populating the `graph` attribute. 

1763 

1764 Notes 

1765 ----- 

1766 This does not read logs, metadata, or packages ; those must always be 

1767 fetched explicitly. 

1768 """ 

1769 self.read_init_quanta() 

1770 self.read_datasets() 

1771 self.read_quanta() 

1772 

1773 def read_datasets(self, datasets: Iterable[uuid.UUID] | None = None) -> None: 

1774 """Read information about the given datasets. 

1775 

1776 Parameters 

1777 ---------- 

1778 datasets : `~collections.abc.Iterable` [`uuid.UUID`], optional 

1779 Iterable of dataset IDs to load. If not provided, all datasets 

1780 will be loaded. The UUIDs and indices of quanta will be ignored. 

1781 """ 

1782 self._read_nodes(datasets, DATASET_ADDRESS_INDEX, DATASET_MB_NAME, ProvenanceDatasetModel) 

1783 

1784 def read_quanta(self, quanta: Iterable[uuid.UUID] | None = None) -> None: 

1785 """Read information about the given quanta. 

1786 

1787 Parameters 

1788 ---------- 

1789 quanta : `~collections.abc.Iterable` [`uuid.UUID`], optional 

1790 Iterable of quantum IDs to load. If not provided, all quanta will 

1791 be loaded. The UUIDs and indices of datasets and special init 

1792 quanta will be ignored. 

1793 """ 

1794 self._read_nodes(quanta, QUANTUM_ADDRESS_INDEX, QUANTUM_MB_NAME, ProvenanceQuantumModel) 

1795 

1796 def _read_nodes( 

1797 self, 

1798 nodes: Iterable[uuid.UUID] | None, 

1799 address_index: int, 

1800 mb_name: str, 

1801 model_type: type[ProvenanceDatasetModel] | type[ProvenanceQuantumModel], 

1802 ) -> None: 

1803 node: ProvenanceDatasetModel | ProvenanceQuantumModel | None 

1804 if nodes is None: 

1805 self.address_reader.read_all() 

1806 nodes = self.address_reader.rows.keys() 

1807 for node in MultiblockReader.read_all_models_in_zip( 

1808 self.zf, 

1809 mb_name, 

1810 model_type, 

1811 self.decompressor, 

1812 int_size=self.header.int_size, 

1813 page_size=self.page_size, 

1814 ): 

1815 if "pipeline_node" in self.graph._bipartite_xgraph.nodes.get(node.node_id, {}): 

1816 # Use the old node to reduce memory usage (since it might 

1817 # also have other outstanding reference holders). 

1818 continue 

1819 node._add_to_graph(self.graph) 

1820 else: 

1821 with MultiblockReader.open_in_zip(self.zf, mb_name, int_size=self.header.int_size) as mb_reader: 

1822 for node_id_or_index in nodes: 1822 ↛ 1823line 1822 didn't jump to line 1823 because the loop on line 1822 never started

1823 address_row = self.address_reader.find(node_id_or_index) 

1824 if "pipeline_node" in self.graph._bipartite_xgraph.nodes.get(address_row.key, {}): 

1825 # Use the old node to reduce memory usage (since it 

1826 # might also have other outstanding reference holders). 

1827 continue 

1828 node = mb_reader.read_model( 

1829 address_row.addresses[address_index], model_type, self.decompressor 

1830 ) 

1831 if node is not None: 

1832 node._add_to_graph(self.graph) 

1833 

1834 def fetch_logs(self, nodes: Iterable[uuid.UUID]) -> dict[uuid.UUID, list[ButlerLogRecords | None]]: 

1835 """Fetch log datasets. 

1836 

1837 Parameters 

1838 ---------- 

1839 nodes : `~collections.abc.Iterable` [ `uuid.UUID` ] 

1840 UUIDs of the log datasets themselves or of the quanta they 

1841 correspond to. 

1842 

1843 Returns 

1844 ------- 

1845 logs : `dict` [ `uuid.UUID`, `list` [\ 

1846 `lsst.daf.butler.ButlerLogRecords` or `None`] ] 

1847 Logs for the given IDs. Each value is a list of 

1848 `lsst.daf.butler.ButlerLogRecords` instances representing different 

1849 execution attempts, ordered chronologically from first to last. 

1850 Attempts where logs were missing will have `None` in this list. 

1851 """ 

1852 result: dict[uuid.UUID, list[ButlerLogRecords | None]] = {} 

1853 with MultiblockReader.open_in_zip(self.zf, LOG_MB_NAME, int_size=self.header.int_size) as mb_reader: 

1854 for node_id_or_index in nodes: 

1855 address_row = self.address_reader.find(node_id_or_index) 

1856 logs_by_attempt = mb_reader.read_model( 

1857 address_row.addresses[LOG_ADDRESS_INDEX], ProvenanceLogRecordsModel, self.decompressor 

1858 ) 

1859 if logs_by_attempt is not None: 1859 ↛ 1854line 1859 didn't jump to line 1854 because the condition on line 1859 was always true

1860 result[node_id_or_index] = [ 

1861 ButlerLogRecords.from_records(attempt_logs) if attempt_logs is not None else None 

1862 for attempt_logs in logs_by_attempt.attempts 

1863 ] 

1864 return result 

1865 

1866 def fetch_metadata(self, nodes: Iterable[uuid.UUID]) -> dict[uuid.UUID, list[TaskMetadata | None]]: 

1867 """Fetch metadata datasets. 

1868 

1869 Parameters 

1870 ---------- 

1871 nodes : `~collections.abc.Iterable` [ `uuid.UUID` ] 

1872 UUIDs of the metadata datasets themselves or of the quanta they 

1873 correspond to. 

1874 

1875 Returns 

1876 ------- 

1877 metadata : `dict` [ `uuid.UUID`, `list` [`.TaskMetadata`] ] 

1878 Metadata for the given IDs. Each value is a list of 

1879 `.TaskMetadata` instances representing different execution 

1880 attempts, ordered chronologically from first to last. Attempts 

1881 where metadata was missing (not written even in the fallback extra 

1882 provenance in the logs) will have `None` in this list. 

1883 """ 

1884 result: dict[uuid.UUID, list[TaskMetadata | None]] = {} 

1885 with MultiblockReader.open_in_zip( 

1886 self.zf, METADATA_MB_NAME, int_size=self.header.int_size 

1887 ) as mb_reader: 

1888 for node_id_or_index in nodes: 

1889 address_row = self.address_reader.find(node_id_or_index) 

1890 metadata_by_attempt = mb_reader.read_model( 

1891 address_row.addresses[METADATA_ADDRESS_INDEX], 

1892 ProvenanceTaskMetadataModel, 

1893 self.decompressor, 

1894 ) 

1895 if metadata_by_attempt is not None: 1895 ↛ 1888line 1895 didn't jump to line 1888 because the condition on line 1895 was always true

1896 result[node_id_or_index] = metadata_by_attempt.attempts 

1897 return result 

1898 

1899 def fetch_packages(self) -> Packages: 

1900 """Fetch package version information.""" 

1901 data = self._read_single_block_raw("packages") 

1902 return Packages.fromBytes(data, format="json") 

1903 

1904 

1905class ProvenanceQuantumGraphWriter: 

1906 """A struct of low-level writer objects for the main components of a 

1907 provenance quantum graph. 

1908 

1909 Parameters 

1910 ---------- 

1911 output_path : `str` 

1912 Path to write the graph to. 

1913 exit_stack : `contextlib.ExitStack` 

1914 Object that can be used to manage multiple context managers. 

1915 log_on_close : `LogOnClose` 

1916 Factory for context managers that log when closed. 

1917 predicted : `.PredictedQuantumGraphComponents` 

1918 Components of the predicted graph. 

1919 zstd_level : `int`, optional 

1920 Compression level. 

1921 cdict_data : `bytes` or `None`, optional 

1922 Bytes representation of the compression dictionary used by the 

1923 compressor. 

1924 loop_wrapper : `~collections.abc.Callable`, optional 

1925 A callable that takes an iterable and returns an equivalent one, to be 

1926 used in all potentially-large loops. This can be used to add progress 

1927 reporting or check for cancelation signals. 

1928 log : `LsstLogAdapter`, optional 

1929 Logger to use for debug messages. 

1930 """ 

1931 

1932 def __init__( 

1933 self, 

1934 output_path: str, 

1935 *, 

1936 exit_stack: ExitStack, 

1937 log_on_close: LogOnClose, 

1938 predicted: PredictedQuantumGraphComponents | PredictedQuantumGraph, 

1939 zstd_level: int = 10, 

1940 cdict_data: bytes | None = None, 

1941 loop_wrapper: LoopWrapper = pass_through, 

1942 log: LsstLogAdapter | None = None, 

1943 ) -> None: 

1944 header = predicted.header.model_copy() 

1945 header.graph_type = "provenance" 

1946 if log is None: 

1947 log = _LOG 

1948 self.log = log 

1949 self._base_writer = exit_stack.enter_context( 

1950 log_on_close.wrap( 

1951 BaseQuantumGraphWriter.open( 

1952 output_path, 

1953 header, 

1954 predicted.pipeline_graph, 

1955 address_filename="nodes", 

1956 zstd_level=zstd_level, 

1957 cdict_data=cdict_data, 

1958 ), 

1959 "Finishing writing provenance quantum graph.", 

1960 ) 

1961 ) 

1962 self._base_writer.address_writer.addresses = [{}, {}, {}, {}] 

1963 self._log_writer = exit_stack.enter_context( 

1964 log_on_close.wrap( 

1965 MultiblockWriter.open_in_zip( 

1966 self._base_writer.zf, LOG_MB_NAME, header.int_size, use_tempfile=True 

1967 ), 

1968 "Copying logs into zip archive.", 

1969 ), 

1970 ) 

1971 self._base_writer.address_writer.addresses[LOG_ADDRESS_INDEX] = self._log_writer.addresses 

1972 self._metadata_writer = exit_stack.enter_context( 

1973 log_on_close.wrap( 

1974 MultiblockWriter.open_in_zip( 

1975 self._base_writer.zf, METADATA_MB_NAME, header.int_size, use_tempfile=True 

1976 ), 

1977 "Copying metadata into zip archive.", 

1978 ) 

1979 ) 

1980 self._base_writer.address_writer.addresses[METADATA_ADDRESS_INDEX] = self._metadata_writer.addresses 

1981 self._dataset_writer = exit_stack.enter_context( 

1982 log_on_close.wrap( 

1983 MultiblockWriter.open_in_zip( 

1984 self._base_writer.zf, DATASET_MB_NAME, header.int_size, use_tempfile=True 

1985 ), 

1986 "Copying dataset provenance into zip archive.", 

1987 ) 

1988 ) 

1989 self._base_writer.address_writer.addresses[DATASET_ADDRESS_INDEX] = self._dataset_writer.addresses 

1990 self._quantum_writer = exit_stack.enter_context( 

1991 log_on_close.wrap( 

1992 MultiblockWriter.open_in_zip( 

1993 self._base_writer.zf, QUANTUM_MB_NAME, header.int_size, use_tempfile=True 

1994 ), 

1995 "Copying quantum provenance into zip archive.", 

1996 ) 

1997 ) 

1998 self._base_writer.address_writer.addresses[QUANTUM_ADDRESS_INDEX] = self._quantum_writer.addresses 

1999 self._init_predicted_quanta(predicted) 

2000 self._populate_xgraph_and_inputs(loop_wrapper) 

2001 self._existing_init_outputs: set[uuid.UUID] = set() 

2002 

2003 def _init_predicted_quanta( 

2004 self, predicted: PredictedQuantumGraph | PredictedQuantumGraphComponents 

2005 ) -> None: 

2006 self._predicted_init_quanta: list[PredictedQuantumDatasetsModel] = [] 

2007 self._predicted_quanta: dict[uuid.UUID, PredictedQuantumDatasetsModel] = {} 

2008 if isinstance(predicted, PredictedQuantumGraph): 

2009 self._predicted_init_quanta.extend(predicted._init_quanta.values()) 

2010 self._predicted_quanta.update(predicted._quantum_datasets) 

2011 else: 

2012 self._predicted_init_quanta.extend(predicted.init_quanta.root) 

2013 self._predicted_quanta.update(predicted.quantum_datasets) 

2014 self._predicted_quanta.update({q.quantum_id: q for q in self._predicted_init_quanta}) 

2015 

2016 def _populate_xgraph_and_inputs(self, loop_wrapper: LoopWrapper = pass_through) -> None: 

2017 self._xgraph = networkx.DiGraph() 

2018 self._overall_inputs: dict[uuid.UUID, PredictedDatasetModel] = {} 

2019 output_dataset_ids: set[uuid.UUID] = set() 

2020 for predicted_quantum in loop_wrapper(self._predicted_quanta.values()): 

2021 if not predicted_quantum.task_label: 

2022 # Skip the 'packages' producer quantum. 

2023 continue 

2024 output_dataset_ids.update(predicted_quantum.iter_output_dataset_ids()) 

2025 for predicted_quantum in loop_wrapper(self._predicted_quanta.values()): 

2026 if not predicted_quantum.task_label: 

2027 # Skip the 'packages' producer quantum. 

2028 continue 

2029 for predicted_input in itertools.chain.from_iterable(predicted_quantum.inputs.values()): 

2030 self._xgraph.add_edge(predicted_input.dataset_id, predicted_quantum.quantum_id) 

2031 if predicted_input.dataset_id not in output_dataset_ids: 

2032 self._overall_inputs.setdefault(predicted_input.dataset_id, predicted_input) 

2033 for predicted_output in itertools.chain.from_iterable(predicted_quantum.outputs.values()): 

2034 self._xgraph.add_edge(predicted_quantum.quantum_id, predicted_output.dataset_id) 

2035 

2036 @property 

2037 def compressor(self) -> Compressor: 

2038 """Object that should be used to compress all JSON blocks.""" 

2039 return self._base_writer.compressor 

2040 

2041 def write_packages(self) -> None: 

2042 """Write package version information to the provenance graph.""" 

2043 packages = Packages.fromSystem(include_all=True) 

2044 data = packages.toBytes("json") 

2045 self._base_writer.write_single_block("packages", data) 

2046 

2047 def write_overall_inputs(self, loop_wrapper: LoopWrapper = pass_through) -> None: 

2048 """Write provenance for overall-input datasets. 

2049 

2050 Parameters 

2051 ---------- 

2052 loop_wrapper : `~collections.abc.Callable`, optional 

2053 A callable that takes an iterable and returns an equivalent one, to 

2054 be used in all potentially-large loops. This can be used to add 

2055 progress reporting or check for cancelation signals. 

2056 """ 

2057 for predicted_input in loop_wrapper(self._overall_inputs.values()): 

2058 if predicted_input.dataset_id not in self._dataset_writer.addresses: 2058 ↛ 2057line 2058 didn't jump to line 2057 because the condition on line 2058 was always true

2059 self._dataset_writer.write_model( 

2060 predicted_input.dataset_id, 

2061 ProvenanceDatasetModel.from_predicted( 

2062 predicted_input, 

2063 producer=None, 

2064 consumers=self._xgraph.successors(predicted_input.dataset_id), 

2065 ), 

2066 self.compressor, 

2067 ) 

2068 del self._overall_inputs 

2069 

2070 def write_init_outputs(self, assume_existence: bool = True) -> None: 

2071 """Write provenance for init-output datasets and init-quanta. 

2072 

2073 Parameters 

2074 ---------- 

2075 assume_existence : `bool`, optional 

2076 If `True`, just assume all init-outputs exist. 

2077 """ 

2078 init_quanta = ProvenanceInitQuantaModel() 

2079 for predicted_init_quantum in self._predicted_init_quanta: 

2080 if not predicted_init_quantum.task_label: 

2081 # Skip the 'packages' producer quantum. 

2082 continue 

2083 for predicted_output in itertools.chain.from_iterable(predicted_init_quantum.outputs.values()): 

2084 provenance_output = ProvenanceDatasetModel.from_predicted( 

2085 predicted_output, 

2086 producer=predicted_init_quantum.quantum_id, 

2087 consumers=self._xgraph.successors(predicted_output.dataset_id), 

2088 ) 

2089 provenance_output.produced = assume_existence or ( 

2090 provenance_output.dataset_id in self._existing_init_outputs 

2091 ) 

2092 self._dataset_writer.write_model( 

2093 provenance_output.dataset_id, provenance_output, self.compressor 

2094 ) 

2095 init_quanta.root.append(ProvenanceInitQuantumModel.from_predicted(predicted_init_quantum)) 

2096 self._base_writer.write_single_model("init_quanta", init_quanta) 

2097 

2098 def write_quantum_provenance( 

2099 self, quantum_id: uuid.UUID, metadata: TaskMetadata | None, logs: ButlerLogRecords | None 

2100 ) -> None: 

2101 """Gather and write provenance for a quantum. 

2102 

2103 Parameters 

2104 ---------- 

2105 quantum_id : `uuid.UUID` 

2106 Unique ID for the quantum. 

2107 metadata : `..TaskMetadata` or `None` 

2108 Task metadata. 

2109 logs : `lsst.daf.butler.logging.ButlerLogRecords` or `None` 

2110 Task logs. 

2111 """ 

2112 predicted_quantum = self._predicted_quanta[quantum_id] 

2113 provenance_models = ProvenanceQuantumScanModels.from_metadata_and_logs( 

2114 predicted_quantum, metadata, logs, incomplete=False 

2115 ) 

2116 scan_data = provenance_models.to_scan_data(predicted_quantum, compressor=self.compressor) 

2117 self.write_scan_data(scan_data) 

2118 

2119 def write_blocked_quantum_provenance(self, quantum_id: uuid.UUID) -> None: 

2120 """Gather and write provenance for a quantum that was blocked by an 

2121 upstream failure. 

2122 

2123 Parameters 

2124 ---------- 

2125 quantum_id : `uuid.UUID` 

2126 Unique ID for the quantum. 

2127 """ 

2128 self.write_scan_data(ProvenanceQuantumScanData.make_blocked(quantum_id)) 

2129 

2130 def write_scan_data(self, scan_data: ProvenanceQuantumScanData) -> None: 

2131 """Write the output of a quantum provenance scan to disk. 

2132 

2133 Parameters 

2134 ---------- 

2135 scan_data : `ProvenanceQuantumScanData` 

2136 Result of a quantum provenance scan. 

2137 """ 

2138 if scan_data.status is ProvenanceQuantumScanStatus.INIT: 

2139 self.log.debug("Handling init-output scan for %s.", scan_data.quantum_id) 

2140 self._existing_init_outputs.update(scan_data.existing_outputs) 

2141 return 

2142 self.log.debug("Handling quantum scan for %s.", scan_data.quantum_id) 

2143 # We shouldn't need this predicted quantum after this method runs; pop 

2144 # from the dict it in the hopes that'll free up some memory when we're 

2145 # done. 

2146 predicted_quantum = self._predicted_quanta.pop(scan_data.quantum_id) 

2147 outputs: dict[uuid.UUID, bytes] = {} 

2148 for predicted_output in itertools.chain.from_iterable(predicted_quantum.outputs.values()): 

2149 provenance_output = ProvenanceDatasetModel.from_predicted( 

2150 predicted_output, 

2151 producer=predicted_quantum.quantum_id, 

2152 consumers=self._xgraph.successors(predicted_output.dataset_id), 

2153 ) 

2154 provenance_output.produced = provenance_output.dataset_id in scan_data.existing_outputs 

2155 outputs[provenance_output.dataset_id] = self.compressor.compress( 

2156 provenance_output.model_dump_json().encode() 

2157 ) 

2158 if not scan_data.quantum: 

2159 scan_data.quantum = ( 

2160 ProvenanceQuantumModel.from_predicted(predicted_quantum).model_dump_json().encode() 

2161 ) 

2162 if scan_data.is_compressed: 

2163 scan_data.quantum = self.compressor.compress(scan_data.quantum) 

2164 if not scan_data.is_compressed: 

2165 scan_data.quantum = self.compressor.compress(scan_data.quantum) 

2166 if scan_data.metadata: 

2167 scan_data.metadata = self.compressor.compress(scan_data.metadata) 

2168 if scan_data.logs: 

2169 scan_data.logs = self.compressor.compress(scan_data.logs) 

2170 self.log.debug("Writing quantum %s.", scan_data.quantum_id) 

2171 self._quantum_writer.write_bytes(scan_data.quantum_id, scan_data.quantum) 

2172 for dataset_id, dataset_data in outputs.items(): 

2173 self._dataset_writer.write_bytes(dataset_id, dataset_data) 

2174 if scan_data.metadata: 

2175 (metadata_output,) = predicted_quantum.outputs[acc.METADATA_OUTPUT_CONNECTION_NAME] 

2176 address = self._metadata_writer.write_bytes(scan_data.quantum_id, scan_data.metadata) 

2177 self._metadata_writer.addresses[metadata_output.dataset_id] = address 

2178 if scan_data.logs: 

2179 (log_output,) = predicted_quantum.outputs[acc.LOG_OUTPUT_CONNECTION_NAME] 

2180 address = self._log_writer.write_bytes(scan_data.quantum_id, scan_data.logs) 

2181 self._log_writer.addresses[log_output.dataset_id] = address 

2182 

2183 

2184class ProvenanceQuantumScanStatus(enum.Enum): 

2185 """Status enum for quantum scanning. 

2186 

2187 Note that this records the status for the *scanning* which is distinct 

2188 from the status of the quantum's execution. 

2189 """ 

2190 

2191 INCOMPLETE = enum.auto() 

2192 """The quantum is not necessarily done running, and cannot be scanned 

2193 conclusively yet. 

2194 """ 

2195 

2196 ABANDONED = enum.auto() 

2197 """The quantum's execution appears to have failed but we cannot rule out 

2198 the possibility that it could be recovered, but we've also waited long 

2199 enough (according to `ScannerTimeConfigDict.retry_timeout`) that it's time 

2200 to stop trying for now. 

2201 

2202 This state means `ProvenanceQuantumScanModels.from_metadata_and_logs` must 

2203 be run again with ``incomplete=False``. 

2204 """ 

2205 

2206 SUCCESSFUL = enum.auto() 

2207 """The quantum was conclusively scanned and was executed successfully, 

2208 unblocking scans for downstream quanta. 

2209 """ 

2210 

2211 FAILED = enum.auto() 

2212 """The quantum was conclusively scanned and failed execution, blocking 

2213 scans for downstream quanta. 

2214 """ 

2215 

2216 BLOCKED = enum.auto() 

2217 """A quantum upstream of this one failed.""" 

2218 

2219 INIT = enum.auto() 

2220 """Init quanta need special handling, because they don't have logs and 

2221 metadata. 

2222 """ 

2223 

2224 

2225@dataclasses.dataclass 

2226class ProvenanceQuantumScanModels: 

2227 """A struct that represents provenance information for a single quantum.""" 

2228 

2229 quantum_id: uuid.UUID 

2230 """Unique ID for the quantum.""" 

2231 

2232 status: ProvenanceQuantumScanStatus = ProvenanceQuantumScanStatus.INCOMPLETE 

2233 """Combined status for the scan and the execution of the quantum.""" 

2234 

2235 attempts: list[ProvenanceQuantumAttemptModel] = dataclasses.field(default_factory=list) 

2236 """Provenance information about each attempt to run the quantum.""" 

2237 

2238 output_existence: dict[uuid.UUID, bool] = dataclasses.field(default_factory=dict) 

2239 """Unique IDs of the output datasets mapped to whether they were actually 

2240 produced. 

2241 """ 

2242 

2243 metadata: ProvenanceTaskMetadataModel = dataclasses.field(default_factory=ProvenanceTaskMetadataModel) 

2244 """Task metadata information for each attempt. 

2245 """ 

2246 

2247 logs: ProvenanceLogRecordsModel = dataclasses.field(default_factory=ProvenanceLogRecordsModel) 

2248 """Log records for each attempt. 

2249 """ 

2250 

2251 @classmethod 

2252 def from_metadata_and_logs( 

2253 cls, 

2254 predicted: PredictedQuantumDatasetsModel, 

2255 metadata: TaskMetadata | None, 

2256 logs: ButlerLogRecords | None, 

2257 *, 

2258 incomplete: bool = False, 

2259 ) -> ProvenanceQuantumScanModels: 

2260 """Construct provenance information from task metadata and logs. 

2261 

2262 Parameters 

2263 ---------- 

2264 predicted : `PredictedQuantumDatasetsModel` 

2265 Information about the predicted quantum. 

2266 metadata : `..TaskMetadata` or `None` 

2267 Task metadata. 

2268 logs : `lsst.daf.butler.logging.ButlerLogRecords` or `None` 

2269 Task logs. 

2270 incomplete : `bool`, optional 

2271 If `True`, treat execution failures as possibly-incomplete quanta 

2272 and do not fully process them; instead just set the status to 

2273 `ProvenanceQuantumScanStatus.ABANDONED` and return. 

2274 

2275 Returns 

2276 ------- 

2277 scan_models : `ProvenanceQuantumScanModels` 

2278 Struct of models that describe quantum provenance. 

2279 

2280 Notes 

2281 ----- 

2282 This method does not necessarily fully populate the `output_existence` 

2283 field; it does what it can given the information in the metadata and 

2284 logs, but the caller is responsible for filling in the existence status 

2285 for any predicted outputs that are not present at all in that `dict`. 

2286 """ 

2287 self = ProvenanceQuantumScanModels(predicted.quantum_id) 

2288 last_attempt = ProvenanceQuantumAttemptModel() 

2289 self._process_logs(predicted, logs, last_attempt, incomplete=incomplete) 

2290 self._process_metadata(predicted, metadata, last_attempt, incomplete=incomplete) 

2291 if self.status is ProvenanceQuantumScanStatus.ABANDONED: 

2292 return self 

2293 self._reconcile_attempts(last_attempt) 

2294 self._extract_output_existence(predicted) 

2295 return self 

2296 

2297 def _process_logs( 

2298 self, 

2299 predicted: PredictedQuantumDatasetsModel, 

2300 logs: ButlerLogRecords | None, 

2301 last_attempt: ProvenanceQuantumAttemptModel, 

2302 *, 

2303 incomplete: bool, 

2304 ) -> None: 

2305 (predicted_log_dataset,) = predicted.outputs[acc.LOG_OUTPUT_CONNECTION_NAME] 

2306 if logs is None: 

2307 self.output_existence[predicted_log_dataset.dataset_id] = False 

2308 if incomplete: 2308 ↛ 2311line 2308 didn't jump to line 2311 because the condition on line 2308 was always true

2309 self.status = ProvenanceQuantumScanStatus.ABANDONED 

2310 else: 

2311 self.status = ProvenanceQuantumScanStatus.FAILED 

2312 else: 

2313 # Set the attempt's run status to FAILED, since the default is 

2314 # UNKNOWN (i.e. logs *and* metadata are missing) and we now know 

2315 # the logs exist. This will usually get replaced by SUCCESSFUL 

2316 # when we look for metadata next. 

2317 last_attempt.status = QuantumAttemptStatus.FAILED 

2318 self.output_existence[predicted_log_dataset.dataset_id] = True 

2319 if logs.extra: 2319 ↛ 2322line 2319 didn't jump to line 2322 because the condition on line 2319 was always true

2320 log_extra = _ExecutionLogRecordsExtra.model_validate(logs.extra) 

2321 self._extract_from_log_extra(log_extra, last_attempt=last_attempt) 

2322 self.logs.attempts.append(list(logs)) 

2323 

2324 def _extract_from_log_extra( 

2325 self, 

2326 log_extra: _ExecutionLogRecordsExtra, 

2327 last_attempt: ProvenanceQuantumAttemptModel | None, 

2328 ) -> None: 

2329 for previous_attempt_log_extra in log_extra.previous_attempts: 

2330 self._extract_from_log_extra( 

2331 previous_attempt_log_extra, 

2332 last_attempt=None, 

2333 ) 

2334 quantum_attempt: ProvenanceQuantumAttemptModel 

2335 if last_attempt is None: 

2336 # This is not the last attempt, so it must be a failure. 

2337 quantum_attempt = ProvenanceQuantumAttemptModel( 

2338 attempt=len(self.attempts), status=QuantumAttemptStatus.FAILED 

2339 ) 

2340 # We also need to get the logs from this extra provenance, since 

2341 # they won't be the main section of the log records. 

2342 self.logs.attempts.append(log_extra.logs) 

2343 # The special last attempt is only appended after we attempt to 

2344 # read metadata later, but we have to append this one now. 

2345 self.attempts.append(quantum_attempt) 

2346 else: 

2347 assert not log_extra.logs, "Logs for the last attempt should not be stored in the extra JSON." 

2348 quantum_attempt = last_attempt 

2349 if log_extra.exception is not None or log_extra.metadata is not None or last_attempt is None: 

2350 # We won't be getting a separate metadata dataset, so anything we 

2351 # might get from the metadata has to come from this extra 

2352 # provenance in the logs. 

2353 quantum_attempt.exception = log_extra.exception 

2354 if log_extra.metadata is not None: 

2355 quantum_attempt.resource_usage = QuantumResourceUsage.from_task_metadata(log_extra.metadata) 

2356 self.metadata.attempts.append(log_extra.metadata) 

2357 else: 

2358 self.metadata.attempts.append(None) 

2359 # Regardless of whether this is the last attempt or not, we can only 

2360 # get the previous_process_quanta from the log extra. 

2361 quantum_attempt.previous_process_quanta.extend(log_extra.previous_process_quanta) 

2362 

2363 def _process_metadata( 

2364 self, 

2365 predicted: PredictedQuantumDatasetsModel, 

2366 metadata: TaskMetadata | None, 

2367 last_attempt: ProvenanceQuantumAttemptModel, 

2368 *, 

2369 incomplete: bool, 

2370 ) -> None: 

2371 (predicted_metadata_dataset,) = predicted.outputs[acc.METADATA_OUTPUT_CONNECTION_NAME] 

2372 if metadata is None: 

2373 self.output_existence[predicted_metadata_dataset.dataset_id] = False 

2374 if incomplete: 

2375 self.status = ProvenanceQuantumScanStatus.ABANDONED 

2376 else: 

2377 self.status = ProvenanceQuantumScanStatus.FAILED 

2378 else: 

2379 self.status = ProvenanceQuantumScanStatus.SUCCESSFUL 

2380 self.output_existence[predicted_metadata_dataset.dataset_id] = True 

2381 last_attempt.status = QuantumAttemptStatus.SUCCESSFUL 

2382 try: 

2383 # Int conversion guards against spurious conversion to 

2384 # float that can apparently sometimes happen in 

2385 # TaskMetadata. 

2386 last_attempt.caveats = QuantumSuccessCaveats(int(metadata["quantum"]["caveats"])) 

2387 except LookupError: 

2388 pass 

2389 try: 

2390 last_attempt.exception = ExceptionInfo._from_metadata( 

2391 metadata[predicted.task_label]["failure"] 

2392 ) 

2393 except LookupError: 

2394 pass 

2395 last_attempt.resource_usage = QuantumResourceUsage.from_task_metadata(metadata) 

2396 self.metadata.attempts.append(metadata) 

2397 

2398 def _reconcile_attempts(self, last_attempt: ProvenanceQuantumAttemptModel) -> None: 

2399 last_attempt.attempt = len(self.attempts) 

2400 self.attempts.append(last_attempt) 

2401 assert self.status is not ProvenanceQuantumScanStatus.INCOMPLETE 

2402 assert self.status is not ProvenanceQuantumScanStatus.ABANDONED 

2403 if len(self.logs.attempts) < len(self.attempts): 2403 ↛ 2407line 2403 didn't jump to line 2407 because the condition on line 2403 was never true

2404 # Logs were not found for this attempt; must have been a hard error 

2405 # that kept the `finally` block from running or otherwise 

2406 # interrupted the writing of the logs. 

2407 self.logs.attempts.append(None) 

2408 if self.status is ProvenanceQuantumScanStatus.SUCCESSFUL: 

2409 # But we found the metadata! Either that hard error happened 

2410 # at a very unlucky time (in between those two writes), or 

2411 # something even weirder happened. 

2412 self.attempts[-1].status = QuantumAttemptStatus.ABORTED_SUCCESS 

2413 else: 

2414 self.attempts[-1].status = QuantumAttemptStatus.FAILED 

2415 if len(self.metadata.attempts) < len(self.attempts): 2415 ↛ 2420line 2415 didn't jump to line 2420 because the condition on line 2415 was never true

2416 # Metadata missing usually just means a failure. In any case, the 

2417 # status will already be correct, either because it was set to a 

2418 # failure when we read the logs, or left at UNKNOWN if there were 

2419 # no logs. Note that scanners never process BLOCKED quanta at all. 

2420 self.metadata.attempts.append(None) 

2421 assert len(self.logs.attempts) == len(self.attempts) or len(self.metadata.attempts) == len( 

2422 self.attempts 

2423 ), ( 

2424 "The only way we can add more than one quantum attempt is by " 

2425 "extracting info stored with the logs, and that always appends " 

2426 "a log attempt and a metadata attempt, so this must be a bug in " 

2427 "this class." 

2428 ) 

2429 

2430 def _extract_output_existence(self, predicted: PredictedQuantumDatasetsModel) -> None: 

2431 try: 

2432 outputs_put = self.metadata.attempts[-1]["quantum"].getArray("outputs") # type: ignore[index] 

2433 except ( 

2434 IndexError, # metadata.attempts is empty 

2435 TypeError, # metadata.attempts[-1] is None 

2436 LookupError, # no 'quantum' entry in metadata or 'outputs' in that 

2437 ): 

2438 pass 

2439 else: 

2440 for id_str in ensure_iterable(outputs_put): 

2441 self.output_existence[uuid.UUID(id_str)] = True 

2442 # If the metadata told us what it wrote, anything not in that 

2443 # list was not written. 

2444 for predicted_output in itertools.chain.from_iterable(predicted.outputs.values()): 

2445 self.output_existence.setdefault(predicted_output.dataset_id, False) 

2446 

2447 def to_scan_data( 

2448 self: ProvenanceQuantumScanModels, 

2449 predicted_quantum: PredictedQuantumDatasetsModel, 

2450 compressor: Compressor | None = None, 

2451 ) -> ProvenanceQuantumScanData: 

2452 """Convert these models to JSON data. 

2453 

2454 Parameters 

2455 ---------- 

2456 predicted_quantum : `PredictedQuantumDatasetsModel` 

2457 Information about the predicted quantum. 

2458 compressor : `Compressor` 

2459 Object that can compress bytes. 

2460 

2461 Returns 

2462 ------- 

2463 scan_data : `ProvenanceQuantumScanData` 

2464 Scan information ready for serialization. 

2465 """ 

2466 quantum: ProvenanceInitQuantumModel | ProvenanceQuantumModel 

2467 if self.status is ProvenanceQuantumScanStatus.INIT: 

2468 quantum = ProvenanceInitQuantumModel.from_predicted(predicted_quantum) 

2469 else: 

2470 quantum = ProvenanceQuantumModel.from_predicted(predicted_quantum) 

2471 quantum.attempts = self.attempts 

2472 for predicted_output in itertools.chain.from_iterable(predicted_quantum.outputs.values()): 

2473 if predicted_output.dataset_id not in self.output_existence: 2473 ↛ 2474line 2473 didn't jump to line 2474 because the condition on line 2473 was never true

2474 raise RuntimeError( 

2475 "Logic bug in provenance gathering or execution invariants: " 

2476 f"no existence information for output {predicted_output.dataset_id} " 

2477 f"({predicted_output.dataset_type_name}@{predicted_output.data_coordinate})." 

2478 ) 

2479 data = ProvenanceQuantumScanData( 

2480 self.quantum_id, 

2481 self.status, 

2482 existing_outputs={ 

2483 dataset_id for dataset_id, was_produced in self.output_existence.items() if was_produced 

2484 }, 

2485 quantum=quantum.model_dump_json().encode(), 

2486 logs=self.logs.model_dump_json().encode() if self.logs.attempts else b"", 

2487 metadata=self.metadata.model_dump_json().encode() if self.metadata.attempts else b"", 

2488 ) 

2489 if compressor is not None: 

2490 data.compress(compressor) 

2491 return data 

2492 

2493 

2494@dataclasses.dataclass 

2495class ProvenanceQuantumScanData: 

2496 """A struct that represents ready-for-serialization provenance information 

2497 for a single quantum. 

2498 """ 

2499 

2500 quantum_id: uuid.UUID 

2501 """Unique ID for the quantum.""" 

2502 

2503 status: ProvenanceQuantumScanStatus 

2504 """Combined status for the scan and the execution of the quantum.""" 

2505 

2506 existing_outputs: set[uuid.UUID] = dataclasses.field(default_factory=set) 

2507 """Unique IDs of the output datasets that were actually written.""" 

2508 

2509 quantum: bytes = b"" 

2510 """Serialized quantum provenance model. 

2511 

2512 This may be empty for quanta that had no attempts. 

2513 """ 

2514 

2515 metadata: bytes = b"" 

2516 """Serialized task metadata.""" 

2517 

2518 logs: bytes = b"" 

2519 """Serialized logs.""" 

2520 

2521 is_compressed: bool = False 

2522 """Whether the ``quantum``, ``metadata``, and ``log`` attributes are 

2523 compressed. 

2524 """ 

2525 

2526 @classmethod 

2527 def make_blocked(cls, quantum_id: uuid.UUID) -> ProvenanceQuantumScanData: 

2528 """Construct provenance information for a quantum blocked by an 

2529 upstream failure. 

2530 

2531 Parameters 

2532 ---------- 

2533 quantum_id : `uuid.UUID` 

2534 Unique ID of the quantum. 

2535 

2536 Returns 

2537 ------- 

2538 scan_data : `ProvenanceQuantumScanData` 

2539 Struct with ready-to-write provenance data. 

2540 """ 

2541 return ProvenanceQuantumScanData( 

2542 quantum_id, 

2543 status=ProvenanceQuantumScanStatus.BLOCKED, 

2544 is_compressed=True, # nothing to compress 

2545 ) 

2546 

2547 def compress(self, compressor: Compressor) -> None: 

2548 """Compress the data in this struct if it has not been compressed 

2549 already. 

2550 

2551 Parameters 

2552 ---------- 

2553 compressor : `Compressor` 

2554 Object with a ``compress`` method that takes and returns `bytes`. 

2555 """ 

2556 if not self.is_compressed: 2556 ↛ exitline 2556 didn't return from function 'compress' because the condition on line 2556 was always true

2557 self.quantum = compressor.compress(self.quantum) 

2558 self.logs = compressor.compress(self.logs) if self.logs else b"" 

2559 self.metadata = compressor.compress(self.metadata) if self.metadata else b"" 

2560 self.is_compressed = True