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-17 21:01 +0000
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-17 21:01 +0000
1# This file is part of pipe_base.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (http://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This software is dual licensed under the GNU General Public License and also
10# under a 3-clause BSD license. Recipients may choose which of these licenses
11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt,
12# respectively. If you choose the GPL option then the following text applies
13# (but note that there is still no warranty even if you opt for BSD instead):
14#
15# This program is free software: you can redistribute it and/or modify
16# it under the terms of the GNU General Public License as published by
17# the Free Software Foundation, either version 3 of the License, or
18# (at your option) any later version.
19#
20# This program is distributed in the hope that it will be useful,
21# but WITHOUT ANY WARRANTY; without even the implied warranty of
22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
23# GNU General Public License for more details.
24#
25# You should have received a copy of the GNU General Public License
26# along with this program. If not, see <http://www.gnu.org/licenses/>.
28from __future__ import annotations
30__all__ = (
31 "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)
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
60import astropy.table
61import networkx
62import numpy as np
63import pydantic
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
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)
99# Sphinx needs imports for type annotations of base class members.
100if "sphinx" in sys.modules:
101 import zipfile # noqa: F401
103 from ._multiblock import AddressReader, Decompressor # noqa: F401
106type LoopWrapper[T] = Callable[[Iterable[T]], Iterable[T]]
108_LOG = getLogger(__file__)
110DATASET_ADDRESS_INDEX = 0
111QUANTUM_ADDRESS_INDEX = 1
112LOG_ADDRESS_INDEX = 2
113METADATA_ADDRESS_INDEX = 3
115DATASET_MB_NAME = "datasets"
116QUANTUM_MB_NAME = "quanta"
117LOG_MB_NAME = "logs"
118METADATA_MB_NAME = "metadata"
121def pass_through[T](arg: T) -> T:
122 return arg
125class ProvenanceDatasetInfo(DatasetInfo):
126 """A typed dictionary that annotates the attributes of the NetworkX graph
127 node data for a provenance dataset.
129 Since NetworkX types are not generic over their node mapping type, this has
130 to be used explicitly, e.g.::
132 node_data: ProvenanceDatasetInfo = xgraph.nodes[dataset_id]
134 where ``xgraph`` is `ProvenanceQuantumGraph.bipartite_xgraph`.
135 """
137 dataset_id: uuid.UUID
138 """Unique identifier for the dataset."""
140 produced: bool
141 """Whether this dataset was produced (vs. only predicted).
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 """
150class ProvenanceQuantumInfo(QuantumInfo):
151 """A typed dictionary that annotates the attributes of the NetworkX graph
152 node data for a provenance quantum.
154 Since NetworkX types are not generic over their node mapping type, this has
155 to be used explicitly, e.g.::
157 node_data: ProvenanceQuantumInfo = xgraph.nodes[quantum_id]
159 where ``xgraph`` is `ProvenanceQuantumGraph.bipartite_xgraph` or
160 `ProvenanceQuantumGraph.quantum_only_xgraph`
161 """
163 status: QuantumAttemptStatus
164 """Enumerated status for the quantum.
166 This corresponds to the last attempt to run this quantum, or
167 `QuantumAttemptStatus.BLOCKED` if there were no attempts.
168 """
170 caveats: QuantumSuccessCaveats | None
171 """Flags indicating caveats on successful quanta.
173 This corresponds to the last attempt to run this quantum.
174 """
176 exception: ExceptionInfo | None
177 """Information about an exception raised when the quantum was executing.
179 This corresponds to the last attempt to run this quantum.
180 """
182 resource_usage: QuantumResourceUsage | None
183 """Resource usage information (timing, memory use) for this quantum.
185 This corresponds to the last attempt to run this quantum.
186 """
188 attempts: list[ProvenanceQuantumAttemptModel]
189 """Information about each attempt to run this quantum.
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 """
196 metadata_id: uuid.UUID
197 """ID of this quantum's metadata dataset."""
199 log_id: uuid.UUID
200 """ID of this quantum's log dataset."""
203class ProvenanceInitQuantumInfo(TypedDict):
204 """A typed dictionary that annotates the attributes of the NetworkX graph
205 node data for a provenance init quantum.
207 Since NetworkX types are not generic over their node mapping type, this has
208 to be used explicitly, e.g.::
210 node_data: ProvenanceInitQuantumInfo = xgraph.nodes[quantum_id]
212 where ``xgraph`` is `ProvenanceQuantumGraph.bipartite_xgraph`.
213 """
215 data_id: DataCoordinate
216 """Data ID of the quantum.
218 This is always an empty ID; this key exists to allow init-quanta and
219 regular quanta to be treated more similarly.
220 """
222 task_label: str
223 """Label of the task for this quantum."""
225 pipeline_node: TaskInitNode
226 """Node in the pipeline graph for this task's init-only step."""
228 config_id: uuid.UUID
229 """ID of this task's config dataset."""
232class ProvenanceDatasetModel(PredictedDatasetModel):
233 """Data model for the datasets in a provenance quantum graph file."""
235 produced: bool
236 """Whether this dataset was produced (vs. only predicted).
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 """
244 producer: uuid.UUID | None = None
245 """ID of the quantum that produced this dataset.
247 This is `None` for overall inputs to the graph.
248 """
250 consumers: list[uuid.UUID] = pydantic.Field(default_factory=list)
251 """IDs of quanta that were predicted to consume this dataset."""
253 @property
254 def node_id(self) -> uuid.UUID:
255 """Alias for the dataset ID."""
256 return self.dataset_id
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.
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.
276 Returns
277 -------
278 provenance : `ProvenanceDatasetModel`
279 Provenance dataset model.
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 )
296 def _add_to_graph(self, graph: ProvenanceQuantumGraph) -> None:
297 """Add this dataset and its edges to quanta to a provenance graph.
299 Parameters
300 ----------
301 graph : `ProvenanceQuantumGraph`
302 Graph to update in place.
304 Notes
305 -----
306 This method adds:
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
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:
337 def copy(self, *args: Any, **kwargs: Any) -> Any:
338 """See `pydantic.BaseModel.copy`."""
339 return super().copy(*args, **kwargs)
341 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
342 """See `pydantic.BaseModel.model_dump`."""
343 return super().model_dump(*args, **kwargs)
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)
349 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
350 """See `pydantic.BaseModel.model_copy`."""
351 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
379class ProvenanceQuantumAttemptModel(pydantic.BaseModel):
380 """Data model for a now-superseded attempt to run a quantum in a
381 provenance quantum graph file.
382 """
384 attempt: int = 0
385 """Counter incremented for every attempt to execute this quantum."""
387 status: QuantumAttemptStatus = QuantumAttemptStatus.UNKNOWN
388 """Enumerated status for the quantum."""
390 caveats: QuantumSuccessCaveats | None = None
391 """Flags indicating caveats on successful quanta."""
393 exception: ExceptionInfo | None = None
394 """Information about an exception raised when the quantum was executing."""
396 resource_usage: QuantumResourceUsage | None = None
397 """Resource usage information (timing, memory use) for this quantum."""
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 """
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:
408 def copy(self, *args: Any, **kwargs: Any) -> Any:
409 """See `pydantic.BaseModel.copy`."""
410 return super().copy(*args, **kwargs)
412 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
413 """See `pydantic.BaseModel.model_dump`."""
414 return super().model_dump(*args, **kwargs)
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)
420 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
421 """See `pydantic.BaseModel.model_copy`."""
422 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
450class ProvenanceLogRecordsModel(pydantic.BaseModel):
451 """Data model for storing execution logs in a provenance quantum graph
452 file.
453 """
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 """
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:
464 def copy(self, *args: Any, **kwargs: Any) -> Any:
465 """See `pydantic.BaseModel.copy`."""
466 return super().copy(*args, **kwargs)
468 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
469 """See `pydantic.BaseModel.model_dump`."""
470 return super().model_dump(*args, **kwargs)
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)
476 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
477 """See `pydantic.BaseModel.model_copy`."""
478 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
506class ProvenanceTaskMetadataModel(pydantic.BaseModel):
507 """Data model for storing task metadata in a provenance quantum graph
508 file.
509 """
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")
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 """
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:
525 def copy(self, *args: Any, **kwargs: Any) -> Any:
526 """See `pydantic.BaseModel.copy`."""
527 return super().copy(*args, **kwargs)
529 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
530 """See `pydantic.BaseModel.model_dump`."""
531 return super().model_dump(*args, **kwargs)
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)
537 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
538 """See `pydantic.BaseModel.model_copy`."""
539 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
567class ProvenanceQuantumReport(pydantic.BaseModel):
568 """A Pydantic model that used to report information about a single
569 (generally problematic) quantum.
570 """
572 quantum_id: uuid.UUID
573 data_id: dict[str, int | str]
574 attempts: list[ProvenanceQuantumAttemptModel]
576 @classmethod
577 def from_info(cls, quantum_id: uuid.UUID, quantum_info: ProvenanceQuantumInfo) -> ProvenanceQuantumReport:
578 """Construct from a provenance quantum graph node.
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 )
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:
597 def copy(self, *args: Any, **kwargs: Any) -> Any:
598 """See `pydantic.BaseModel.copy`."""
599 return super().copy(*args, **kwargs)
601 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
602 """See `pydantic.BaseModel.model_dump`."""
603 return super().model_dump(*args, **kwargs)
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)
609 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
610 """See `pydantic.BaseModel.model_copy`."""
611 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
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 """
644 root: dict[TaskLabel, dict[str, dict[str | None, list[ProvenanceQuantumReport]]]] = {}
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:
650 def copy(self, *args: Any, **kwargs: Any) -> Any:
651 """See `pydantic.BaseModel.copy`."""
652 return super().copy(*args, **kwargs)
654 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
655 """See `pydantic.BaseModel.model_dump`."""
656 return super().model_dump(*args, **kwargs)
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)
662 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
663 """See `pydantic.BaseModel.model_copy`."""
664 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
692class ProvenanceQuantumModel(pydantic.BaseModel):
693 """Data model for the quanta in a provenance quantum graph file."""
695 quantum_id: uuid.UUID
696 """Unique identifier for the quantum."""
698 task_label: TaskLabel
699 """Name of the type of this dataset."""
701 data_coordinate: DataCoordinateValues = pydantic.Field(default_factory=list)
702 """The full values (required and implied) of this dataset's data ID."""
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 """
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 """
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.
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 """
723 @property
724 def node_id(self) -> uuid.UUID:
725 """Alias for the quantum ID."""
726 return self.quantum_id
728 @classmethod
729 def from_predicted(cls, predicted: PredictedQuantumDatasetsModel) -> ProvenanceQuantumModel:
730 """Construct from a predicted quantum model.
732 Parameters
733 ----------
734 predicted : `PredictedQuantumDatasetsModel`
735 Information about the quantum from the predicted graph.
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 )
758 def _add_to_graph(self, graph: ProvenanceQuantumGraph) -> None:
759 """Add this quantum and its edges to datasets to a provenance graph.
761 Parameters
762 ----------
763 graph : `ProvenanceQuantumGraph`
764 Graph to update in place.
766 Notes
767 -----
768 This method adds:
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)
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:
849 def copy(self, *args: Any, **kwargs: Any) -> Any:
850 """See `pydantic.BaseModel.copy`."""
851 return super().copy(*args, **kwargs)
853 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
854 """See `pydantic.BaseModel.model_dump`."""
855 return super().model_dump(*args, **kwargs)
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)
861 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
862 """See `pydantic.BaseModel.model_copy`."""
863 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
891class ProvenanceInitQuantumModel(pydantic.BaseModel):
892 """Data model for the special "init" quanta in a provenance quantum graph
893 file.
894 """
896 quantum_id: uuid.UUID
897 """Unique identifier for the quantum."""
899 task_label: TaskLabel
900 """Name of the type of this dataset.
902 This is always a parent dataset type name, not a component.
904 Note that full dataset type definitions are stored in the pipeline graph.
905 """
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 """
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 """
917 @classmethod
918 def from_predicted(cls, predicted: PredictedQuantumDatasetsModel) -> ProvenanceInitQuantumModel:
919 """Construct from a predicted quantum model.
921 Parameters
922 ----------
923 predicted : `PredictedQuantumDatasetsModel`
924 Information about the quantum from the predicted graph.
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 )
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.
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.
956 Notes
957 -----
958 This method adds:
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
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:
1000 def copy(self, *args: Any, **kwargs: Any) -> Any:
1001 """See `pydantic.BaseModel.copy`."""
1002 return super().copy(*args, **kwargs)
1004 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
1005 """See `pydantic.BaseModel.model_dump`."""
1006 return super().model_dump(*args, **kwargs)
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)
1012 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
1013 """See `pydantic.BaseModel.model_copy`."""
1014 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
1042class ProvenanceInitQuantaModel(pydantic.RootModel):
1043 """Data model for the init quanta in a provenance graph."""
1045 root: list[ProvenanceInitQuantumModel] = pydantic.Field(default_factory=list)
1046 """List of special "init" quanta, one for each task."""
1048 def _add_to_graph(self, graph: ProvenanceQuantumGraph) -> None:
1049 """Add this quantum and its edges to datasets to a provenance graph.
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)
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:
1064 def copy(self, *args: Any, **kwargs: Any) -> Any:
1065 """See `pydantic.BaseModel.copy`."""
1066 return super().copy(*args, **kwargs)
1068 def model_dump(self, *args: Any, **kwargs: Any) -> Any:
1069 """See `pydantic.BaseModel.model_dump`."""
1070 return super().model_dump(*args, **kwargs)
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)
1076 def model_copy(self, *args: Any, **kwargs: Any) -> Any:
1077 """See `pydantic.BaseModel.model_copy`."""
1078 return super().model_copy(*args, **kwargs)
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)
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)
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)
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)
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)
1106class ProvenanceQuantumGraph(BaseQuantumGraph):
1107 """A quantum graph that represents processing that has already been
1108 executed.
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.
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 """
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 }
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.
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.
1168 Returns
1169 -------
1170 context : `contextlib.AbstractContextManager`
1171 A context manager that yields a tuple of
1173 - the `ProvenanceQuantumGraph`
1174 - the `Butler` constructed (or `None`)
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
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.
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.::
1224 info: ProvenanceInitQuantumInfo = qg.bipartite_xgraph.nodes[id]
1225 """
1226 return self._init_quanta
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.
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).
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
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.
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).
1255 Reading a quantum also populates its log and metadata datasets.
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
1262 @property
1263 def quantum_only_xgraph(self) -> networkx.DiGraph:
1264 """A directed acyclic graph with quanta as nodes (and datasets elided).
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.
1276 Node attributes are described by the `ProvenanceQuantumInfo` types.
1278 This graph does not include special "init" quanta.
1280 The returned object is a read-only view of an internal one.
1281 """
1282 return self._quantum_only_xgraph.copy(as_view=True)
1284 @property
1285 def bipartite_xgraph(self) -> networkx.DiGraph:
1286 """A directed acyclic graph with quantum and dataset nodes.
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.
1298 Node attributes are described by the
1299 `ProvenanceQuantumInfo`, `ProvenanceInitQuantumInfo`, and
1300 `ProvenanceDatasetInfo` types.
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).
1307 The returned object is a read-only view of an internal one.
1308 """
1309 return self._bipartite_xgraph.copy(as_view=True)
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.
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.
1323 expand_caveats : `bool`, optional
1324 Whether to display a comma-separated list of task caveats instead
1325 of a collapsed `multiple` marker.
1327 as_table : `bool`, optional
1328 Whether to return an `astropy.table.Table` or a list of rows.
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.
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)
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
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.
1396 Parameters
1397 ----------
1398 as_table : `bool`, optional
1399 Whether to return an `astropy.table.Table` or a list of rows.
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
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.
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.
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())
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.
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).
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
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.
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"
1635 if output_format == "json":
1636 json_obj: dict[str, Any] = {"tasks": None, "exceptions": None}
1637 expand_caveats = True
1639 if TYPE_CHECKING:
1640 assert json_obj
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
1686 if output_format == "json":
1687 json_obj["legend"] = QuantumSuccessCaveats.legend()
1688 json.dump(json_obj, sys.stdout)
1691@dataclasses.dataclass
1692class ProvenanceQuantumGraphReader(BaseQuantumGraphReader):
1693 """A helper class for reading provenance quantum graphs.
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`.
1701 The various ``read_*`` methods in this class update the `graph` attribute
1702 in place.
1703 """
1705 graph: ProvenanceQuantumGraph = dataclasses.field(init=False)
1706 """Loaded provenance graph, populated in place as components are read."""
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.
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.
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
1748 def __post_init__(self) -> None:
1749 self.graph = ProvenanceQuantumGraph(self.header, self.pipeline_graph)
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)
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.
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()
1773 def read_datasets(self, datasets: Iterable[uuid.UUID] | None = None) -> None:
1774 """Read information about the given datasets.
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)
1784 def read_quanta(self, quanta: Iterable[uuid.UUID] | None = None) -> None:
1785 """Read information about the given quanta.
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)
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)
1834 def fetch_logs(self, nodes: Iterable[uuid.UUID]) -> dict[uuid.UUID, list[ButlerLogRecords | None]]:
1835 """Fetch log datasets.
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.
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
1866 def fetch_metadata(self, nodes: Iterable[uuid.UUID]) -> dict[uuid.UUID, list[TaskMetadata | None]]:
1867 """Fetch metadata datasets.
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.
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
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")
1905class ProvenanceQuantumGraphWriter:
1906 """A struct of low-level writer objects for the main components of a
1907 provenance quantum graph.
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 """
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()
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})
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)
2036 @property
2037 def compressor(self) -> Compressor:
2038 """Object that should be used to compress all JSON blocks."""
2039 return self._base_writer.compressor
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)
2047 def write_overall_inputs(self, loop_wrapper: LoopWrapper = pass_through) -> None:
2048 """Write provenance for overall-input datasets.
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
2070 def write_init_outputs(self, assume_existence: bool = True) -> None:
2071 """Write provenance for init-output datasets and init-quanta.
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)
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.
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)
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.
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))
2130 def write_scan_data(self, scan_data: ProvenanceQuantumScanData) -> None:
2131 """Write the output of a quantum provenance scan to disk.
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
2184class ProvenanceQuantumScanStatus(enum.Enum):
2185 """Status enum for quantum scanning.
2187 Note that this records the status for the *scanning* which is distinct
2188 from the status of the quantum's execution.
2189 """
2191 INCOMPLETE = enum.auto()
2192 """The quantum is not necessarily done running, and cannot be scanned
2193 conclusively yet.
2194 """
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.
2202 This state means `ProvenanceQuantumScanModels.from_metadata_and_logs` must
2203 be run again with ``incomplete=False``.
2204 """
2206 SUCCESSFUL = enum.auto()
2207 """The quantum was conclusively scanned and was executed successfully,
2208 unblocking scans for downstream quanta.
2209 """
2211 FAILED = enum.auto()
2212 """The quantum was conclusively scanned and failed execution, blocking
2213 scans for downstream quanta.
2214 """
2216 BLOCKED = enum.auto()
2217 """A quantum upstream of this one failed."""
2219 INIT = enum.auto()
2220 """Init quanta need special handling, because they don't have logs and
2221 metadata.
2222 """
2225@dataclasses.dataclass
2226class ProvenanceQuantumScanModels:
2227 """A struct that represents provenance information for a single quantum."""
2229 quantum_id: uuid.UUID
2230 """Unique ID for the quantum."""
2232 status: ProvenanceQuantumScanStatus = ProvenanceQuantumScanStatus.INCOMPLETE
2233 """Combined status for the scan and the execution of the quantum."""
2235 attempts: list[ProvenanceQuantumAttemptModel] = dataclasses.field(default_factory=list)
2236 """Provenance information about each attempt to run the quantum."""
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 """
2243 metadata: ProvenanceTaskMetadataModel = dataclasses.field(default_factory=ProvenanceTaskMetadataModel)
2244 """Task metadata information for each attempt.
2245 """
2247 logs: ProvenanceLogRecordsModel = dataclasses.field(default_factory=ProvenanceLogRecordsModel)
2248 """Log records for each attempt.
2249 """
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.
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.
2275 Returns
2276 -------
2277 scan_models : `ProvenanceQuantumScanModels`
2278 Struct of models that describe quantum provenance.
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
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))
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)
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)
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 )
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)
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.
2454 Parameters
2455 ----------
2456 predicted_quantum : `PredictedQuantumDatasetsModel`
2457 Information about the predicted quantum.
2458 compressor : `Compressor`
2459 Object that can compress bytes.
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
2494@dataclasses.dataclass
2495class ProvenanceQuantumScanData:
2496 """A struct that represents ready-for-serialization provenance information
2497 for a single quantum.
2498 """
2500 quantum_id: uuid.UUID
2501 """Unique ID for the quantum."""
2503 status: ProvenanceQuantumScanStatus
2504 """Combined status for the scan and the execution of the quantum."""
2506 existing_outputs: set[uuid.UUID] = dataclasses.field(default_factory=set)
2507 """Unique IDs of the output datasets that were actually written."""
2509 quantum: bytes = b""
2510 """Serialized quantum provenance model.
2512 This may be empty for quanta that had no attempts.
2513 """
2515 metadata: bytes = b""
2516 """Serialized task metadata."""
2518 logs: bytes = b""
2519 """Serialized logs."""
2521 is_compressed: bool = False
2522 """Whether the ``quantum``, ``metadata``, and ``log`` attributes are
2523 compressed.
2524 """
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.
2531 Parameters
2532 ----------
2533 quantum_id : `uuid.UUID`
2534 Unique ID of the quantum.
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 )
2547 def compress(self, compressor: Compressor) -> None:
2548 """Compress the data in this struct if it has not been compressed
2549 already.
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