Source code for pycif.plugins.obsoperators.standard.transforms.period_pipe

from logging import info

from .utils.default_subsimus import default_subsimus
from .utils.fwd_pipe import init_pipe_links, compute_pipe_order


[docs] def period_pipe(self, all_transforms, mapper): """Arrange all transforms into ordered forward and adjoint execution pipes. Determines the chronologically correct execution order for every ``(transform, sub-simulation date)`` pair by: 1. Propagating sub-simulation periods from each transform to its precursors and successors via :func:`~.utils.default_subsimus.default_subsimus`. 2. Building the graph of links between transforms in forward and adjoint directions with :func:`~.utils.fwd_pipe.init_pipe_links`. 3. Without Dask, walking these graphs to get the sequential execution order with :func:`~.utils.fwd_pipe.compute_pipe_order`. Each returned pipe is a list of ``(date, transform_id, direction)`` tuples, where ``direction`` is either ``'forward'`` or ``'adjoint'`` and controls whether a transform runs in its normal or dry-run mode. Args: self (ObsOperator): the obs-operator plugin instance. all_transforms: the :class:`~pycif.utils.classes.transforms.Transform` object holding all initialized transforms. mapper (dict): the pipeline mapper dictionary mapping transform IDs to their sub-simulation, input/output and precursor/successor metadata. Returns: tuple[list, list]: ``(pipe_fwd, pipe_adj)`` where each element is a list of ``(datetime.datetime, str, str)`` tuples giving the execution order for forward and adjoint runs respectively. """ info("Computing the optimal order of transformation. This can take a while") # First update subsimulations from precursors and successors default_subsimus(all_transforms, mapper) for mode in ["forward", "adjoint"]: info(f"Doing {mode} order") pipe_links, transforms_ids = init_pipe_links( self, all_transforms, mapper, mode=mode) # Dask handles the order itself, but batch computation # needs the forward order if self.use_dask and not ( mode == "forward" and hasattr(self, "batch_computation")): continue compute_pipe_order( self, all_transforms, mapper, pipe_links, transforms_ids, mode=mode)