Skip to content

omnipy.compute.flow

Flow definitions for composing tasks and subflows.

This module exposes Omnipy's public flow types. Use LinearFlowTemplate for sequential pipelines, DagFlowTemplate for dependency-driven directed acyclic graphs, and FuncFlowTemplate for flows expressed as a single coordinating callable.

ATTRIBUTE DESCRIPTION
LinearFlowTemplate

Decorator-style template factory for sequential flows.

TYPE: Any

DagFlowTemplate

Decorator-style template factory for directed acyclic graph flows.

TYPE: Any

FuncFlowTemplate

Decorator-style template factory for callable-backed coordinating flows.

TYPE: Any

CLASS DESCRIPTION
DagFlow

Execute a flow whose child jobs exchange data through named DAG keys.

DagFlowTemplateCore
FlowBase

Provide a shared marker base for Omnipy flow objects.

FuncFlow

Execute a flow backed by one coordinating callable.

FuncFlowTemplateCore

Implement the core template behavior for function flows.

LinearFlow

Execute a flow whose tasks run in declaration order.

LinearFlowTemplateCore

Implement the core template behavior for linear flows.

FUNCTION DESCRIPTION
DagFlowTemplate

Decorator-style factory for defining directed acyclic graph flows.

FuncFlowTemplate

Decorator-style factory for defining callable-backed coordinating flows.

LinearFlowTemplate

Decorator-style factory for defining sequential flows.

DagFlow

Bases: JobMixin[IsDagFlowTemplate[_CallP, _RetT], IsDagFlow[_CallP, _RetT], _CallP, _RetT], FlowBase, ChildJobListArgJobBase[IsDagFlowTemplate[_CallP, _RetT], IsDagFlow[_CallP, _RetT], _CallP, _RetT], Generic[_CallP, _RetT]


              flowchart BT
              omnipy.compute.flow.DagFlow[DagFlow]
              omnipy.compute._job.JobMixin[JobMixin]
              omnipy.compute.flow.FlowBase[FlowBase]
              omnipy.compute._joblist_job.ChildJobListArgJobBase[ChildJobListArgJobBase]
              omnipy.compute._func_job.FuncArgJobBase[FuncArgJobBase]
              omnipy.compute._func_job.PlainFuncArgJobBase[PlainFuncArgJobBase]
              omnipy.compute._job.JobBase[JobBase]
              omnipy.hub.log.mixin.LogMixin[LogMixin]
              omnipy.util.mixin.DynamicMixinAcceptor[DynamicMixinAcceptor]

                              omnipy.compute._job.JobMixin --> omnipy.compute.flow.DagFlow
                                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobMixin
                

                omnipy.compute.flow.FlowBase --> omnipy.compute.flow.DagFlow
                
                omnipy.compute._joblist_job.ChildJobListArgJobBase --> omnipy.compute.flow.DagFlow
                                omnipy.compute._func_job.FuncArgJobBase --> omnipy.compute._joblist_job.ChildJobListArgJobBase
                                omnipy.compute._func_job.PlainFuncArgJobBase --> omnipy.compute._func_job.FuncArgJobBase
                                omnipy.compute._job.JobBase --> omnipy.compute._func_job.PlainFuncArgJobBase
                                omnipy.hub.log.mixin.LogMixin --> omnipy.compute._job.JobBase
                
                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobBase
                






              click omnipy.compute.flow.DagFlow href "" "omnipy.compute.flow.DagFlow"
              click omnipy.compute._job.JobMixin href "" "omnipy.compute._job.JobMixin"
              click omnipy.compute.flow.FlowBase href "" "omnipy.compute.flow.FlowBase"
              click omnipy.compute._joblist_job.ChildJobListArgJobBase href "" "omnipy.compute._joblist_job.ChildJobListArgJobBase"
              click omnipy.compute._func_job.FuncArgJobBase href "" "omnipy.compute._func_job.FuncArgJobBase"
              click omnipy.compute._func_job.PlainFuncArgJobBase href "" "omnipy.compute._func_job.PlainFuncArgJobBase"
              click omnipy.compute._job.JobBase href "" "omnipy.compute._job.JobBase"
              click omnipy.hub.log.mixin.LogMixin href "" "omnipy.hub.log.mixin.LogMixin"
              click omnipy.util.mixin.DynamicMixinAcceptor href "" "omnipy.util.mixin.DynamicMixinAcceptor"
            

Execute a flow whose child jobs exchange data through named DAG keys.

A DagFlow routes values between child jobs by accumulated keyword name rather than by one strict positional chain. Use it for branching and joining pipelines whose dependencies form a directed acyclic graph.

Instances are typically produced by calling a DagFlowTemplate rather than by constructing DagFlow directly.

CLASS DESCRIPTION
DataClassAndJobParentInfo
METHOD DESCRIPTION
__init__
accept_mixin

Register a mixin class for dynamic composition.

create_job

Create an applied job instance from the concrete job class.

log

Emit a log message, optionally using an explicit event timestamp.

reset_mixins

Clear all accepted mixins and restore the original init signature.

revise

Return a template reconstructed from this applied job.

ATTRIBUTE DESCRIPTION
callable_type

TYPE: CallableType.Literals

child_job_templates

TYPE: tuple[ChildJobTemplateLike, ...]

config

Return the job configuration visible to this instance.

TYPE: IsJobConfig

engine

Return the engine associated with this job, if any.

TYPE: IsEngine | None

in_flow_context

Return whether the job is currently executing inside a flow context.

TYPE: bool

logger

Return the logger bound to the concrete instance type.

TYPE: Logger

time_of_cur_toplevel_flow_run

Return the start time of the active top-level flow run, if any.

TYPE: datetime | None

Source code in src/omnipy/compute/flow.py
class DagFlow(
        JobMixin[IsDagFlowTemplate[_CallP, _RetT], IsDagFlow[_CallP, _RetT], _CallP, _RetT],
        FlowBase,
        ChildJobListArgJobBase[
            IsDagFlowTemplate[_CallP, _RetT],
            IsDagFlow[_CallP, _RetT],
            _CallP,
            _RetT,
        ],
        Generic[_CallP, _RetT],
):
    """Execute a flow whose child jobs exchange data through named DAG keys.

    A ``DagFlow`` routes values between child jobs by accumulated keyword name
    rather than by one strict positional chain. Use it for branching and
    joining pipelines whose dependencies form a directed acyclic graph.

    Instances are typically produced by calling a ``DagFlowTemplate`` rather
    than by constructing ``DagFlow`` directly.
    """
    def _apply_engine_decorator(self, engine: IsEngine) -> None:
        """Register the engine decorator for DAG-flow execution.

        When a runner engine is bound to the flow, this method asks that
        engine to wrap the flow's call function with DAG-specific execution
        behavior.

        Args:
            self: Current DAG flow instance.
            engine: Engine candidate supplied during job setup.
        """
        if self.engine:
            engine = cast(IsJobRunnerEngine, self.engine)
            self_with_mixins = cast(IsDagFlow, self)
            engine.apply_job_decorator(
                JobType.DAG_FLOW,
                self_with_mixins,
                self._accept_call_func_decorator,
            )

    @classmethod
    def _get_job_template_subcls_for_revise(cls) -> type[IsDagFlowTemplate[_CallP, _RetT]]:
        """Return the template type used to revise this flow.

        Revision operations use this hook to recover the decorator-backed DAG
        flow template class corresponding to an executable DAG flow instance.

        Returns:
            type[IsDagFlowTemplate[_CallP, _RetT]]: The [DagFlowTemplate][]
                class associated with this flow.
        """
        return cast(type[IsDagFlowTemplate[_CallP, _RetT]], DagFlowTemplateCore)

callable_type property

callable_type: CallableType.Literals

child_job_templates property

child_job_templates: tuple[ChildJobTemplateLike, ...]

config property

config: IsJobConfig

Return the job configuration visible to this instance.

RETURNS DESCRIPTION
IsJobConfig

Active job configuration used for runtime behavior.

TYPE: IsJobConfig

engine property

engine: IsEngine | None

Return the engine associated with this job, if any.

RETURNS DESCRIPTION
IsEngine | None

IsEngine | None: Engine used for decoration and execution, or None.

in_flow_context property

in_flow_context: bool

Return whether the job is currently executing inside a flow context.

RETURNS DESCRIPTION
bool

True when a surrounding flow context is active.

TYPE: bool

logger property

logger: Logger

Return the logger bound to the concrete instance type.

RETURNS DESCRIPTION
Logger

Logger used by the object for Omnipy log messages.

TYPE: Logger

time_of_cur_toplevel_flow_run property

time_of_cur_toplevel_flow_run: datetime | None

Return the start time of the active top-level flow run, if any.

RETURNS DESCRIPTION
datetime | None

datetime | None: Timestamp for the current outermost flow run, or None.

DataClassAndJobParentInfo

Bases: NamedTuple


              flowchart BT
              omnipy.compute.flow.DagFlow.DataClassAndJobParentInfo[DataClassAndJobParentInfo]

              

              click omnipy.compute.flow.DagFlow.DataClassAndJobParentInfo href "" "omnipy.compute.flow.DagFlow.DataClassAndJobParentInfo"
            
ATTRIBUTE DESCRIPTION
dataset_or_model

TYPE: Literal['dataset', 'model']

parent_coerces_from_kwargs

TYPE: bool

Source code in src/omnipy/compute/_joblist_job.py
class DataClassAndJobParentInfo(NamedTuple):
    dataset_or_model: Literal['dataset', 'model']
    parent_coerces_from_kwargs: bool

dataset_or_model instance-attribute

dataset_or_model: Literal['dataset', 'model']

parent_coerces_from_kwargs instance-attribute

parent_coerces_from_kwargs: bool

__init__

__init__(*args, **kwargs)
Source code in src/omnipy/compute/_job.py
def __init__(self, *args, **kwargs):
    if JobBase not in self.__class__.__mro__:
        raise TypeError('JobMixin is not meant to be instantiated outside the context '
                        'of a JobBase subclass.')

accept_mixin classmethod

accept_mixin(mixin_cls: Type) -> None

Register a mixin class for dynamic composition.

PARAMETER DESCRIPTION
mixin_cls

Mixin class whose __init__ keyword-only parameters should be merged into the acceptor signature.

TYPE: Type

Source code in src/omnipy/util/mixin.py
@classmethod
def accept_mixin(cls, mixin_cls: Type) -> None:
    """Register a mixin class for dynamic composition.

    Args:
        mixin_cls: Mixin class whose ``__init__`` keyword-only parameters
            should be merged into the acceptor signature.
    """
    cls._accept_mixin(mixin_cls, update=True)

create_job classmethod

create_job(*args: object, **kwargs: object) -> _JobT

Create an applied job instance from the concrete job class.

PARAMETER DESCRIPTION
*args

Positional constructor arguments.

TYPE: object DEFAULT: ()

**kwargs

Keyword constructor arguments.

TYPE: object DEFAULT: {}

RETURNS DESCRIPTION
_JobT

New applied job instance.

TYPE: _JobT

Source code in src/omnipy/compute/_job.py
@classmethod
def create_job(cls, *args: object, **kwargs: object) -> _JobT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOB_CREATE_JOB_SUMMARY}}
    #
    # {{ISJOB_CREATE_JOB_DETAILS}}
    """Create an applied job instance from the concrete job class.

    Args:
        *args: Positional constructor arguments.
        **kwargs: Keyword constructor arguments.

    Returns:
        _JobT: New applied job instance.
    """
    cls_as_job_base = cast(IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT], cls)
    return cls_as_job_base._create_job(*args, **kwargs)

log

log(log_msg: str, level: int = INFO, datetime_obj: datetime | None = None)

Emit a log message, optionally using an explicit event timestamp.

PARAMETER DESCRIPTION
log_msg

Message text to send to the logger.

TYPE: str

level

Standard library logging level.

TYPE: int DEFAULT: INFO

datetime_obj

Timestamp to attach to the record instead of wall-clock time.

TYPE: datetime | None DEFAULT: None

Source code in src/omnipy/hub/log/mixin.py
def log(self, log_msg: str, level: int = INFO, datetime_obj: datetime | None = None):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{CANLOG_LOG_SUMMARY}}
    #
    # {{CANLOG_LOG_DETAILS}}
    """Emit a log message, optionally using an explicit event timestamp.

    Args:
        log_msg: Message text to send to the logger.
        level: Standard library logging level.
        datetime_obj: Timestamp to attach to the record instead of wall-clock time.
    """
    if self._logger is not None:
        create_time = datetime_obj.timestamp() if datetime_obj else time.time()
        self._logger.log(level, log_msg, extra=dict(timestamp=create_time))

reset_mixins classmethod

reset_mixins()

Clear all accepted mixins and restore the original init signature.

Source code in src/omnipy/util/mixin.py
@classmethod
def reset_mixins(cls):
    """Clear all accepted mixins and restore the original init signature."""
    cls._mixin_classes.clear()
    cls._init_params_per_mixin_cls.clear()
    cls.__init__.__signature__ = cls._orig_init_signature

revise

revise() -> _JobTemplateT

Return a template reconstructed from this applied job.

RETURNS DESCRIPTION
_JobTemplateT

Template carrying the current job configuration.

TYPE: _JobTemplateT

Source code in src/omnipy/compute/_job.py
def revise(self) -> _JobTemplateT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOB_REVISE_SUMMARY}}
    #
    # {{ISJOB_REVISE_DETAILS}}
    """Return a template reconstructed from this applied job.

    Returns:
        _JobTemplateT: Template carrying the current job configuration.
    """
    self_as_job_base = cast(
        IsJobBase[IsJobTemplate[_JobTemplateT, _JobT, _CallP, _RetT], _JobT, _CallP, _RetT],
        self)
    job_template = self_as_job_base._revise()
    update_wrapper(job_template, self, updated=[])
    return cast(_JobTemplateT, job_template)

DagFlowTemplateCore

Bases: ChildJobListArgJobBase[IsDagFlowTemplate[_CallP, _RetT], IsDagFlow[_CallP, _RetT], _CallP, _RetT], JobTemplateMixin[IsDagFlowTemplate[_CallP, _RetT], IsDagFlow[_CallP, _RetT], _CallP, _RetT], FlowBase, Generic[_CallP, _RetT]


              flowchart BT
              omnipy.compute.flow.DagFlowTemplateCore[DagFlowTemplateCore]
              omnipy.compute._joblist_job.ChildJobListArgJobBase[ChildJobListArgJobBase]
              omnipy.compute._func_job.FuncArgJobBase[FuncArgJobBase]
              omnipy.compute._func_job.PlainFuncArgJobBase[PlainFuncArgJobBase]
              omnipy.compute._job.JobBase[JobBase]
              omnipy.hub.log.mixin.LogMixin[LogMixin]
              omnipy.util.mixin.DynamicMixinAcceptor[DynamicMixinAcceptor]
              omnipy.compute._job.JobTemplateMixin[JobTemplateMixin]
              omnipy.compute.flow.FlowBase[FlowBase]

                              omnipy.compute._joblist_job.ChildJobListArgJobBase --> omnipy.compute.flow.DagFlowTemplateCore
                                omnipy.compute._func_job.FuncArgJobBase --> omnipy.compute._joblist_job.ChildJobListArgJobBase
                                omnipy.compute._func_job.PlainFuncArgJobBase --> omnipy.compute._func_job.FuncArgJobBase
                                omnipy.compute._job.JobBase --> omnipy.compute._func_job.PlainFuncArgJobBase
                                omnipy.hub.log.mixin.LogMixin --> omnipy.compute._job.JobBase
                
                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobBase
                




                omnipy.compute._job.JobTemplateMixin --> omnipy.compute.flow.DagFlowTemplateCore
                
                omnipy.compute.flow.FlowBase --> omnipy.compute.flow.DagFlowTemplateCore
                


              click omnipy.compute.flow.DagFlowTemplateCore href "" "omnipy.compute.flow.DagFlowTemplateCore"
              click omnipy.compute._joblist_job.ChildJobListArgJobBase href "" "omnipy.compute._joblist_job.ChildJobListArgJobBase"
              click omnipy.compute._func_job.FuncArgJobBase href "" "omnipy.compute._func_job.FuncArgJobBase"
              click omnipy.compute._func_job.PlainFuncArgJobBase href "" "omnipy.compute._func_job.PlainFuncArgJobBase"
              click omnipy.compute._job.JobBase href "" "omnipy.compute._job.JobBase"
              click omnipy.hub.log.mixin.LogMixin href "" "omnipy.hub.log.mixin.LogMixin"
              click omnipy.util.mixin.DynamicMixinAcceptor href "" "omnipy.util.mixin.DynamicMixinAcceptor"
              click omnipy.compute._job.JobTemplateMixin href "" "omnipy.compute._job.JobTemplateMixin"
              click omnipy.compute.flow.FlowBase href "" "omnipy.compute.flow.FlowBase"
            
CLASS DESCRIPTION
DataClassAndJobParentInfo
METHOD DESCRIPTION
__init__
accept_mixin

Register a mixin class for dynamic composition.

apply

Create an applied job from this template without executing it.

create_job_template

Create a job template instance from the concrete template class.

log

Emit a log message, optionally using an explicit event timestamp.

refine

Forward refinement to the shared template lifecycle implementation.

reset_mixins

Clear all accepted mixins and restore the original init signature.

run

Apply the template and execute the resulting job immediately.

ATTRIBUTE DESCRIPTION
callable_type

TYPE: CallableType.Literals

child_job_templates

TYPE: tuple[ChildJobTemplateLike, ...]

config

Return the job configuration visible to this instance.

TYPE: IsJobConfig

engine

Return the engine associated with this job, if any.

TYPE: IsEngine | None

in_flow_context

Return whether the job is currently executing inside a flow context.

TYPE: bool

logger

Return the logger bound to the concrete instance type.

TYPE: Logger

Source code in src/omnipy/compute/flow.py
 942
 943
 944
 945
 946
 947
 948
 949
 950
 951
 952
 953
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
class DagFlowTemplateCore(ChildJobListArgJobBase[IsDagFlowTemplate[_CallP, _RetT],
                                                 IsDagFlow[_CallP, _RetT],
                                                 _CallP,
                                                 _RetT],
                          JobTemplateMixin[IsDagFlowTemplate[_CallP, _RetT],
                                           IsDagFlow[_CallP, _RetT],
                                           _CallP,
                                           _RetT],
                          FlowBase,
                          Generic[_CallP, _RetT]):
    _coerce_data_class_children_from_kwargs: ClassVar[bool] = True

    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # Implement the core template behavior for DAG flows.
    #
    # {{DAG_FLOW_TEMPLATE_DESCRIPTION}}
    #
    # Instances are normally produced through the [DagFlowTemplate][] decorator
    # factory rather than by direct construction.
    #
    """Implement the core template behavior for DAG flows.

    A DAG flow template wraps a Python callable together with child job
    templates whose dependencies form a directed acyclic graph. Use this
    when flow steps branch and join but must not form cycles.

    ### Decorator usage

    Apply the template factory as a decorator to a Python callable.
    The wrapped callable becomes a reusable job template whose public outer
    signature is visible to template users and to the applied jobs created
    from it.

    ### Outer callable and child jobs

    The wrapped callable defines the public outer signature of the flow,
    while the child-job list defines the executed body.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True)
        ... def uppercase(data_file: TextModel) -> TextModel:
        ...     return data_file.content.upper()

        >>> @om.TaskTemplate()
        ... def join_texts(
        ...     upper: TextDataset,
        ...     original: TextDataset,
        ... ) -> TextDataset:
        ...     merged = TextDataset()
        ...     for title in upper:
        ...         merged[title] = f'{upper[title].content}|{original[title].content}'
        ...     return merged

        >>> @om.DagFlowTemplate(
        ...     uppercase.refine(result_key='upper'),
        ...     join_texts.refine(param_key_map={'upper': 'upper', 'original': 'dataset'}),
        ... )
        ... def my_dag(
        ...     dataset: TextDataset,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'HI|hi', 'b': 'BYE|bye'})
        >>> my_dag.run(text_files) == expected
        True

    ### Child jobs and data classes

    Child-job templates may be
    [TaskTemplate][omnipy.compute.task.TaskTemplate] or flow-template
    instances.
    This lets flows nest other flows as well as terminal tasks.

    In linear and DAG flows, Model and Dataset subclasses are also allowed
    as child-job entries. During apply, they are coerced into helper task
    templates that construct the requested data object from positional
    input in linear flows or from keyword-matched input in DAG flows.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.LinearFlowTemplate(TextModel)
        ... def wrap_text(raw_text: str) -> TextModel:
        ...     return TextModel(raw_text)

        >>> wrap_text.run('hello').content
        'hello'
        >>> @om.DagFlowTemplate(TextDataset)
        ... def collect_texts(first: str, second: str) -> TextDataset:
        ...     return TextDataset({'first': first, 'second': second})

        >>> collect_texts.run(first='hello', second='bye') == TextDataset({
        ...     'first': 'hello',
        ...     'second': 'bye',
        ... })
        True

    ### Callable-type validation

    For linear and DAG flows, the outer callable is primarily declarative:
    its signature exposes the public flow interface, while child jobs define
    the executed body.

    The outer callable type is validated against the child-job composition.
    The terminal child determines whether the flow behaves like a function
    or generator, and any async child lifts the full flow to an async
    callable type.

    As a result, a sync-function outer callable fits sync child execution,
    a generator outer callable fits generator-producing terminal children,
    and async outer callables are required when child composition is async.
    Mismatches raise ``TypeError`` when the flow template is created.

    ### Generator shorthand with ``Void``

    When the validated outer callable must be a generator or
    async-generator only to expose the correct public signature, use
    [Void][omnipy.compute.helpers.Void] in the body:

    Examples:
        >>> import omnipy as om
        >>> from collections.abc import Iterator

        >>> @om.TaskTemplate()
        ... def emit_lines() -> Iterator[str]:
        ...     yield 'first'
        ...     yield 'second'

        >>> @om.LinearFlowTemplate(emit_lines)
        ... def line_stream() -> Iterator[str]:
        ...     yield from om.Void()

    This shorthand exists only to satisfy the declared outer callable type;
    the child jobs still perform the actual flow work.

    ### Outer signature and modifiers

    The wrapped callable's parameter list and return annotation define the
    outer interface of the template.

    ``fixed_params`` permanently supplies selected callable parameters.

    ``param_key_map`` renames selected callable parameters to external
    keyword names that callers or parent flows use when supplying inputs.

    ``iterate_over_data_files``, ``output_dataset_param``, and
    ``output_dataset_cls`` adapt that outer interface for dataset-wise
    iteration.

    When ``iterate_over_data_files=True`` and the inner first parameter is
    annotated as ``Model[T]``, callers see an outer
    ``dataset: Dataset[Model[T]]`` parameter and the outer return type
    becomes a dataset of the per-item return type. The inner callable still
    receives one model object at a time.

    ``result_key`` wraps the returned value in a single-key dictionary,
    which is especially useful when a downstream DAG step should receive
    the result under a predictable name.

    Examples:
        >>> # With modifiers
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_other(number: int, other: int) -> int:
        ...     return number + other

        >>> plus_one = plus_other.refine(fixed_params={'other': 1})
        >>> plus_one.run(4)
        5
        >>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
        >>> plus_x.run(4, x=3)
        7
        >>> plus_one_dict = plus_one.refine(result_key='number')
        >>> plus_one_dict.run(4)
        {'number': 5}

    Examples:
        >>> # With dataset-wise iteration
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> add_suffix.run(text_files, suffix='!') == expected
        True

    ### DAG orchestration

    DAG flows route values by keyword name instead of chaining every child
    result positionally.

    The outer flow call first binds its arguments to the outer callable
    signature. Each child then receives the matching keyword subset from the
    accumulated named results.

    By default, a non-dictionary child result is stored under the child job
    name. ``result_key`` stores it under a custom key instead, which is the
    usual way to make one branch feed another. Dictionary results merge
    directly into the accumulated named result set.

    ``param_key_map`` lets a child read from externally visible DAG keys
    using different internal callable parameter names, and ``fixed_params``
    pins selected child inputs regardless of what earlier branches produce.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate()
        ... def uppercase(data_file: TextModel) -> TextModel:
        ...     return data_file.content.upper()

        >>> @om.TaskTemplate()
        ... def append_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> @om.TaskTemplate()
        ... def combine_texts(
        ...     left_dataset: TextDataset,
        ...     right_dataset: TextDataset,
        ... ) -> TextDataset:
        ...     merged = TextDataset()
        ...     for title in left_dataset:
        ...         merged[title] = (
        ...             f'{left_dataset[title].content}|'
        ...             f'{right_dataset[title].content}'
        ...         )
        ...     return merged

        >>> @om.DagFlowTemplate(
        ...     uppercase.refine(
        ...         iterate_over_data_files=True,
        ...         result_key='upper',
        ...     ),
        ...     append_suffix.refine(
        ...         iterate_over_data_files=True,
        ...         result_key='suffixed',
        ...         param_key_map={'suffix': 'ending'},
        ...     ),
        ...     combine_texts.refine(
        ...         param_key_map={
        ...             'left_dataset': 'upper',
        ...             'right_dataset': 'suffixed',
        ...         },
        ...     ),
        ... )
        ... def dag_flow(
        ...     dataset: TextDataset,
        ...     ending: str,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({
        ...     'a': 'HI|hi!',
        ...     'b': 'BYE|bye!',
        ... })
        >>> dag_flow.run(text_files, ending='!') == expected
        True

    ### Tasks and flows

    Tasks are terminal jobs: they wrap one callable and execute one compute
    step.

    Flows are orchestration jobs: they may contain child tasks and child
    flows, so larger pipelines can be assembled hierarchically from smaller
    reusable pieces.

    ### Lifecycle

    Apply a template with [`apply()`][omnipy.compute._job.JobTemplateMixin.apply]
    to create a runnable job with engine decorators and current config attached.
    Call the resulting applied job with runtime arguments.

    Use [`run()`][omnipy.compute._job.JobTemplateMixin.run] as a shorthand for
    ``apply()`` followed immediately by calling the applied job.

    Use [`refine()`][omnipy.compute._job.JobTemplateMixin.refine] to reuse a
    template while changing configuration such as ``name``, ``fixed_params``,
    or ``param_key_map``.

    Use [`revise()`][omnipy.compute._job.JobMixin.revise] on an applied job to
    reconstruct a template from that job's current configuration.

    Examples:
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_one(number: int) -> int:
        ...     return number + 1

        >>> plus_one.run(1)
        2
        >>> applied_job = plus_one.apply()
        >>> applied_job(2)
        3
        >>> refined_template = plus_one.refine(name='plus_one_renamed')
        >>> revised_template = applied_job.revise()

    Instances are normally produced through the [DagFlowTemplate][] decorator
    factory rather than by direct construction.
    """
    @classmethod
    def _get_job_subcls_for_apply(cls) -> type[IsDagFlow[_CallP, _RetT]]:
        """Return the executable DAG flow type produced by this template.

        The template/application machinery calls this hook when it needs the
        concrete flow class to instantiate from a DAG flow template.

        Returns:
            type[IsDagFlow[_CallP, _RetT]]: The executable [DagFlow][] subclass
                associated with this template.
        """
        return cast(type[IsDagFlow[_CallP, _RetT]], DagFlow[_CallP, _RetT])

callable_type property

callable_type: CallableType.Literals

child_job_templates property

child_job_templates: tuple[ChildJobTemplateLike, ...]

config property

config: IsJobConfig

Return the job configuration visible to this instance.

RETURNS DESCRIPTION
IsJobConfig

Active job configuration used for runtime behavior.

TYPE: IsJobConfig

engine property

engine: IsEngine | None

Return the engine associated with this job, if any.

RETURNS DESCRIPTION
IsEngine | None

IsEngine | None: Engine used for decoration and execution, or None.

in_flow_context property

in_flow_context: bool

Return whether the job is currently executing inside a flow context.

RETURNS DESCRIPTION
bool

True when a surrounding flow context is active.

TYPE: bool

logger property

logger: Logger

Return the logger bound to the concrete instance type.

RETURNS DESCRIPTION
Logger

Logger used by the object for Omnipy log messages.

TYPE: Logger

DataClassAndJobParentInfo

Bases: NamedTuple


              flowchart BT
              omnipy.compute.flow.DagFlowTemplateCore.DataClassAndJobParentInfo[DataClassAndJobParentInfo]

              

              click omnipy.compute.flow.DagFlowTemplateCore.DataClassAndJobParentInfo href "" "omnipy.compute.flow.DagFlowTemplateCore.DataClassAndJobParentInfo"
            
ATTRIBUTE DESCRIPTION
dataset_or_model

TYPE: Literal['dataset', 'model']

parent_coerces_from_kwargs

TYPE: bool

Source code in src/omnipy/compute/_joblist_job.py
class DataClassAndJobParentInfo(NamedTuple):
    dataset_or_model: Literal['dataset', 'model']
    parent_coerces_from_kwargs: bool

dataset_or_model instance-attribute

dataset_or_model: Literal['dataset', 'model']

parent_coerces_from_kwargs instance-attribute

parent_coerces_from_kwargs: bool

__init__

__init__(
    job_func: Callable[_CallP, _RetT],
    /,
    *child_job_templates: ChildJobTemplateLike,
    **kwargs: object,
) -> None
Source code in src/omnipy/compute/_joblist_job.py
def __init__(self,
             job_func: Callable[_CallP, _RetT],
             /,
             *child_job_templates: ChildJobTemplateLike,
             **kwargs: object) -> None:
    super().__init__(job_func, *child_job_templates, **kwargs)
    self._child_job_templates: tuple[ChildJobTemplateLike, ...] = child_job_templates
    self._validate_callable_type_against_child_job_templates()

accept_mixin classmethod

accept_mixin(mixin_cls: Type) -> None

Register a mixin class for dynamic composition.

PARAMETER DESCRIPTION
mixin_cls

Mixin class whose __init__ keyword-only parameters should be merged into the acceptor signature.

TYPE: Type

Source code in src/omnipy/util/mixin.py
@classmethod
def accept_mixin(cls, mixin_cls: Type) -> None:
    """Register a mixin class for dynamic composition.

    Args:
        mixin_cls: Mixin class whose ``__init__`` keyword-only parameters
            should be merged into the acceptor signature.
    """
    cls._accept_mixin(mixin_cls, update=True)

apply

apply() -> _JobT

Create an applied job from this template without executing it.

RETURNS DESCRIPTION
_JobT

Applied job instance ready to be called.

TYPE: _JobT

Source code in src/omnipy/compute/_job.py
def apply(self) -> _JobT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_APPLY_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_APPLY_DETAILS}}
    """Create an applied job from this template without executing it.

    Returns:
        _JobT: Applied job instance ready to be called.
    """
    job = self._cast_to_job_tmpl()._apply()
    update_wrapper(job, self, updated=[])
    return cast(_JobT, job)

create_job_template classmethod

create_job_template(*args: object, **kwargs: object) -> _JobTemplateT

Create a job template instance from the concrete template class.

PARAMETER DESCRIPTION
*args

Positional constructor arguments.

TYPE: object DEFAULT: ()

**kwargs

Keyword constructor arguments.

TYPE: object DEFAULT: {}

RETURNS DESCRIPTION
_JobTemplateT

New job template instance.

TYPE: _JobTemplateT

Source code in src/omnipy/compute/_job.py
@classmethod
def create_job_template(cls, *args: object, **kwargs: object) -> _JobTemplateT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_CREATE_JOB_TEMPLATE_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_CREATE_JOB_TEMPLATE_DETAILS}}
    """Create a job template instance from the concrete template class.

    Args:
        *args: Positional constructor arguments.
        **kwargs: Keyword constructor arguments.

    Returns:
        _JobTemplateT: New job template instance.
    """

    cls_as_job_base = cast(type[IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT]], cls)
    return cls_as_job_base._create_job_template(*args, **kwargs)

log

log(log_msg: str, level: int = INFO, datetime_obj: datetime | None = None)

Emit a log message, optionally using an explicit event timestamp.

PARAMETER DESCRIPTION
log_msg

Message text to send to the logger.

TYPE: str

level

Standard library logging level.

TYPE: int DEFAULT: INFO

datetime_obj

Timestamp to attach to the record instead of wall-clock time.

TYPE: datetime | None DEFAULT: None

Source code in src/omnipy/hub/log/mixin.py
def log(self, log_msg: str, level: int = INFO, datetime_obj: datetime | None = None):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{CANLOG_LOG_SUMMARY}}
    #
    # {{CANLOG_LOG_DETAILS}}
    """Emit a log message, optionally using an explicit event timestamp.

    Args:
        log_msg: Message text to send to the logger.
        level: Standard library logging level.
        datetime_obj: Timestamp to attach to the record instead of wall-clock time.
    """
    if self._logger is not None:
        create_time = datetime_obj.timestamp() if datetime_obj else time.time()
        self._logger.log(level, log_msg, extra=dict(timestamp=create_time))

refine

refine(*args: Any, update: bool = True, **kwargs: object) -> _JobTemplateT

Forward refinement to the shared template lifecycle implementation.

See IsFuncArgJobTemplate.refine and IsChildJobListArgJobTemplate.refine.

Source code in src/omnipy/compute/_job.py
def refine(self, *args: Any, update: bool = True, **kwargs: object) -> _JobTemplateT:
    """Forward refinement to the shared template lifecycle implementation.

    See [`IsFuncArgJobTemplate.refine`]
    [omnipy.shared.protocols.compute.job.IsFuncArgJobTemplate.refine] and
    [`IsChildJobListArgJobTemplate.refine`]
    [omnipy.shared.protocols.compute.job.IsChildJobListArgJobTemplate.refine].
    """
    self_as_job_base = cast(IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT], self)
    return self_as_job_base._refine(*args, update=update, **kwargs)

reset_mixins classmethod

reset_mixins()

Clear all accepted mixins and restore the original init signature.

Source code in src/omnipy/util/mixin.py
@classmethod
def reset_mixins(cls):
    """Clear all accepted mixins and restore the original init signature."""
    cls._mixin_classes.clear()
    cls._init_params_per_mixin_cls.clear()
    cls.__init__.__signature__ = cls._orig_init_signature

run

run(*args: _CallP.args, **kwargs: _CallP.kwargs) -> _RetT

Apply the template and execute the resulting job immediately.

PARAMETER DESCRIPTION
*args

Positional arguments passed to the applied job.

TYPE: _CallP.args DEFAULT: ()

**kwargs

Keyword arguments passed to the applied job.

TYPE: _CallP.kwargs DEFAULT: {}

RETURNS DESCRIPTION
_RetCovT

Result returned by the applied job.

TYPE: _RetT

Source code in src/omnipy/compute/_job.py
def run(self, *args: _CallP.args, **kwargs: _CallP.kwargs) -> _RetT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_RUN_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_RUN_DETAILS}}
    """Apply the template and execute the resulting job immediately.

    Args:
        *args: Positional arguments passed to the applied job.
        **kwargs: Keyword arguments passed to the applied job.

    Returns:
        _RetCovT: Result returned by the applied job.
    """
    # TODO: Using JobTemplateMixin.run() inside flows should give error message
    return self._cast_to_job_tmpl().apply()(*args, **kwargs)

FlowBase

Provide a shared marker base for Omnipy flow objects.

FlowBase exists to give concrete flow templates and executable flow instances a common nominal base type. It does not define behavior on its own, but it makes flow-specific mixin registration and type-based checks possible within the compute subsystem.

Source code in src/omnipy/compute/flow.py
class FlowBase:
    """Provide a shared marker base for Omnipy flow objects.

    ``FlowBase`` exists to give concrete flow templates and executable flow
    instances a common nominal base type. It does not define behavior on its
    own, but it makes flow-specific mixin registration and type-based checks
    possible within the compute subsystem.
    """

    ...

FuncFlow

Bases: JobMixin[IsFuncFlowTemplate[_CallP, _RetT], IsFuncFlow[_CallP, _RetT], _CallP, _RetT], FlowBase, FuncArgJobBase[IsFuncFlowTemplate[_CallP, _RetT], IsFuncFlow[_CallP, _RetT], _CallP, _RetT], Generic[_CallP, _RetT]


              flowchart BT
              omnipy.compute.flow.FuncFlow[FuncFlow]
              omnipy.compute._job.JobMixin[JobMixin]
              omnipy.compute.flow.FlowBase[FlowBase]
              omnipy.compute._func_job.FuncArgJobBase[FuncArgJobBase]
              omnipy.compute._func_job.PlainFuncArgJobBase[PlainFuncArgJobBase]
              omnipy.compute._job.JobBase[JobBase]
              omnipy.hub.log.mixin.LogMixin[LogMixin]
              omnipy.util.mixin.DynamicMixinAcceptor[DynamicMixinAcceptor]

                              omnipy.compute._job.JobMixin --> omnipy.compute.flow.FuncFlow
                                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobMixin
                

                omnipy.compute.flow.FlowBase --> omnipy.compute.flow.FuncFlow
                
                omnipy.compute._func_job.FuncArgJobBase --> omnipy.compute.flow.FuncFlow
                                omnipy.compute._func_job.PlainFuncArgJobBase --> omnipy.compute._func_job.FuncArgJobBase
                                omnipy.compute._job.JobBase --> omnipy.compute._func_job.PlainFuncArgJobBase
                                omnipy.hub.log.mixin.LogMixin --> omnipy.compute._job.JobBase
                
                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobBase
                





              click omnipy.compute.flow.FuncFlow href "" "omnipy.compute.flow.FuncFlow"
              click omnipy.compute._job.JobMixin href "" "omnipy.compute._job.JobMixin"
              click omnipy.compute.flow.FlowBase href "" "omnipy.compute.flow.FlowBase"
              click omnipy.compute._func_job.FuncArgJobBase href "" "omnipy.compute._func_job.FuncArgJobBase"
              click omnipy.compute._func_job.PlainFuncArgJobBase href "" "omnipy.compute._func_job.PlainFuncArgJobBase"
              click omnipy.compute._job.JobBase href "" "omnipy.compute._job.JobBase"
              click omnipy.hub.log.mixin.LogMixin href "" "omnipy.hub.log.mixin.LogMixin"
              click omnipy.util.mixin.DynamicMixinAcceptor href "" "omnipy.util.mixin.DynamicMixinAcceptor"
            

Execute a flow backed by one coordinating callable.

A FuncFlow runs the wrapped callable itself as the flow body instead of orchestrating an explicit child-job list. Use it when one callable already captures the desired control flow and runtime behavior.

Instances are typically produced by calling a FuncFlowTemplate rather than by constructing FuncFlow directly.

METHOD DESCRIPTION
__init__
accept_mixin

Register a mixin class for dynamic composition.

create_job

Create an applied job instance from the concrete job class.

log

Emit a log message, optionally using an explicit event timestamp.

reset_mixins

Clear all accepted mixins and restore the original init signature.

revise

Return a template reconstructed from this applied job.

ATTRIBUTE DESCRIPTION
callable_type

TYPE: CallableType.Literals

config

Return the job configuration visible to this instance.

TYPE: IsJobConfig

engine

Return the engine associated with this job, if any.

TYPE: IsEngine | None

in_flow_context

Return whether the job is currently executing inside a flow context.

TYPE: bool

logger

Return the logger bound to the concrete instance type.

TYPE: Logger

time_of_cur_toplevel_flow_run

Return the start time of the active top-level flow run, if any.

TYPE: datetime | None

Source code in src/omnipy/compute/flow.py
class FuncFlow(
        JobMixin[
            IsFuncFlowTemplate[_CallP, _RetT],
            IsFuncFlow[_CallP, _RetT],
            _CallP,
            _RetT,
        ],
        FlowBase,
        FuncArgJobBase[
            IsFuncFlowTemplate[_CallP, _RetT],
            IsFuncFlow[_CallP, _RetT],
            _CallP,
            _RetT,
        ],
        Generic[_CallP, _RetT],
):
    """Execute a flow backed by one coordinating callable.

    A ``FuncFlow`` runs the wrapped callable itself as the flow body instead of
    orchestrating an explicit child-job list. Use it when one callable already
    captures the desired control flow and runtime behavior.

    Instances are typically produced by calling a ``FuncFlowTemplate`` rather
    than by constructing ``FuncFlow`` directly.
    """
    def _apply_engine_decorator(self, engine: IsEngine) -> None:
        """Register the engine decorator for callable-backed flow execution.

        When a runner engine is bound to the flow, this method asks that
        engine to wrap the coordinating callable with function-flow execution
        behavior.

        Args:
            self: Current function flow instance.
            engine: Engine candidate supplied during job setup.
        """
        if self.engine:
            engine = cast(IsJobRunnerEngine, self.engine)
            self_with_mixins = cast(IsFuncFlow, self)
            engine.apply_job_decorator(
                JobType.FUNC_FLOW,
                self_with_mixins,
                self._accept_call_func_decorator,
            )

    @classmethod
    def _get_job_template_subcls_for_revise(cls) -> type[IsFuncFlowTemplate[_CallP, _RetT]]:
        """Return the template type used to revise this flow.

        Revision operations use this hook to recover the decorator-backed
        function flow template class corresponding to an executable function
        flow instance.

        Returns:
            type[IsFuncFlowTemplate[_CallP, _RetT]]: The [FuncFlowTemplate][]
                class associated with this flow.
        """
        return cast(type[IsFuncFlowTemplate[_CallP, _RetT]], FuncFlowTemplateCore)

callable_type property

callable_type: CallableType.Literals

config property

config: IsJobConfig

Return the job configuration visible to this instance.

RETURNS DESCRIPTION
IsJobConfig

Active job configuration used for runtime behavior.

TYPE: IsJobConfig

engine property

engine: IsEngine | None

Return the engine associated with this job, if any.

RETURNS DESCRIPTION
IsEngine | None

IsEngine | None: Engine used for decoration and execution, or None.

in_flow_context property

in_flow_context: bool

Return whether the job is currently executing inside a flow context.

RETURNS DESCRIPTION
bool

True when a surrounding flow context is active.

TYPE: bool

logger property

logger: Logger

Return the logger bound to the concrete instance type.

RETURNS DESCRIPTION
Logger

Logger used by the object for Omnipy log messages.

TYPE: Logger

time_of_cur_toplevel_flow_run property

time_of_cur_toplevel_flow_run: datetime | None

Return the start time of the active top-level flow run, if any.

RETURNS DESCRIPTION
datetime | None

datetime | None: Timestamp for the current outermost flow run, or None.

__init__

__init__(*args, **kwargs)
Source code in src/omnipy/compute/_job.py
def __init__(self, *args, **kwargs):
    if JobBase not in self.__class__.__mro__:
        raise TypeError('JobMixin is not meant to be instantiated outside the context '
                        'of a JobBase subclass.')

accept_mixin classmethod

accept_mixin(mixin_cls: Type) -> None

Register a mixin class for dynamic composition.

PARAMETER DESCRIPTION
mixin_cls

Mixin class whose __init__ keyword-only parameters should be merged into the acceptor signature.

TYPE: Type

Source code in src/omnipy/util/mixin.py
@classmethod
def accept_mixin(cls, mixin_cls: Type) -> None:
    """Register a mixin class for dynamic composition.

    Args:
        mixin_cls: Mixin class whose ``__init__`` keyword-only parameters
            should be merged into the acceptor signature.
    """
    cls._accept_mixin(mixin_cls, update=True)

create_job classmethod

create_job(*args: object, **kwargs: object) -> _JobT

Create an applied job instance from the concrete job class.

PARAMETER DESCRIPTION
*args

Positional constructor arguments.

TYPE: object DEFAULT: ()

**kwargs

Keyword constructor arguments.

TYPE: object DEFAULT: {}

RETURNS DESCRIPTION
_JobT

New applied job instance.

TYPE: _JobT

Source code in src/omnipy/compute/_job.py
@classmethod
def create_job(cls, *args: object, **kwargs: object) -> _JobT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOB_CREATE_JOB_SUMMARY}}
    #
    # {{ISJOB_CREATE_JOB_DETAILS}}
    """Create an applied job instance from the concrete job class.

    Args:
        *args: Positional constructor arguments.
        **kwargs: Keyword constructor arguments.

    Returns:
        _JobT: New applied job instance.
    """
    cls_as_job_base = cast(IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT], cls)
    return cls_as_job_base._create_job(*args, **kwargs)

log

log(log_msg: str, level: int = INFO, datetime_obj: datetime | None = None)

Emit a log message, optionally using an explicit event timestamp.

PARAMETER DESCRIPTION
log_msg

Message text to send to the logger.

TYPE: str

level

Standard library logging level.

TYPE: int DEFAULT: INFO

datetime_obj

Timestamp to attach to the record instead of wall-clock time.

TYPE: datetime | None DEFAULT: None

Source code in src/omnipy/hub/log/mixin.py
def log(self, log_msg: str, level: int = INFO, datetime_obj: datetime | None = None):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{CANLOG_LOG_SUMMARY}}
    #
    # {{CANLOG_LOG_DETAILS}}
    """Emit a log message, optionally using an explicit event timestamp.

    Args:
        log_msg: Message text to send to the logger.
        level: Standard library logging level.
        datetime_obj: Timestamp to attach to the record instead of wall-clock time.
    """
    if self._logger is not None:
        create_time = datetime_obj.timestamp() if datetime_obj else time.time()
        self._logger.log(level, log_msg, extra=dict(timestamp=create_time))

reset_mixins classmethod

reset_mixins()

Clear all accepted mixins and restore the original init signature.

Source code in src/omnipy/util/mixin.py
@classmethod
def reset_mixins(cls):
    """Clear all accepted mixins and restore the original init signature."""
    cls._mixin_classes.clear()
    cls._init_params_per_mixin_cls.clear()
    cls.__init__.__signature__ = cls._orig_init_signature

revise

revise() -> _JobTemplateT

Return a template reconstructed from this applied job.

RETURNS DESCRIPTION
_JobTemplateT

Template carrying the current job configuration.

TYPE: _JobTemplateT

Source code in src/omnipy/compute/_job.py
def revise(self) -> _JobTemplateT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOB_REVISE_SUMMARY}}
    #
    # {{ISJOB_REVISE_DETAILS}}
    """Return a template reconstructed from this applied job.

    Returns:
        _JobTemplateT: Template carrying the current job configuration.
    """
    self_as_job_base = cast(
        IsJobBase[IsJobTemplate[_JobTemplateT, _JobT, _CallP, _RetT], _JobT, _CallP, _RetT],
        self)
    job_template = self_as_job_base._revise()
    update_wrapper(job_template, self, updated=[])
    return cast(_JobTemplateT, job_template)

FuncFlowTemplateCore

Bases: FuncArgJobBase[IsFuncFlowTemplate[_CallP, _RetT], IsFuncFlow[_CallP, _RetT], _CallP, _RetT], JobTemplateMixin[IsFuncFlowTemplate[_CallP, _RetT], IsFuncFlow[_CallP, _RetT], _CallP, _RetT], FlowBase, Generic[_CallP, _RetT]


              flowchart BT
              omnipy.compute.flow.FuncFlowTemplateCore[FuncFlowTemplateCore]
              omnipy.compute._func_job.FuncArgJobBase[FuncArgJobBase]
              omnipy.compute._func_job.PlainFuncArgJobBase[PlainFuncArgJobBase]
              omnipy.compute._job.JobBase[JobBase]
              omnipy.hub.log.mixin.LogMixin[LogMixin]
              omnipy.util.mixin.DynamicMixinAcceptor[DynamicMixinAcceptor]
              omnipy.compute._job.JobTemplateMixin[JobTemplateMixin]
              omnipy.compute.flow.FlowBase[FlowBase]

                              omnipy.compute._func_job.FuncArgJobBase --> omnipy.compute.flow.FuncFlowTemplateCore
                                omnipy.compute._func_job.PlainFuncArgJobBase --> omnipy.compute._func_job.FuncArgJobBase
                                omnipy.compute._job.JobBase --> omnipy.compute._func_job.PlainFuncArgJobBase
                                omnipy.hub.log.mixin.LogMixin --> omnipy.compute._job.JobBase
                
                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobBase
                



                omnipy.compute._job.JobTemplateMixin --> omnipy.compute.flow.FuncFlowTemplateCore
                
                omnipy.compute.flow.FlowBase --> omnipy.compute.flow.FuncFlowTemplateCore
                


              click omnipy.compute.flow.FuncFlowTemplateCore href "" "omnipy.compute.flow.FuncFlowTemplateCore"
              click omnipy.compute._func_job.FuncArgJobBase href "" "omnipy.compute._func_job.FuncArgJobBase"
              click omnipy.compute._func_job.PlainFuncArgJobBase href "" "omnipy.compute._func_job.PlainFuncArgJobBase"
              click omnipy.compute._job.JobBase href "" "omnipy.compute._job.JobBase"
              click omnipy.hub.log.mixin.LogMixin href "" "omnipy.hub.log.mixin.LogMixin"
              click omnipy.util.mixin.DynamicMixinAcceptor href "" "omnipy.util.mixin.DynamicMixinAcceptor"
              click omnipy.compute._job.JobTemplateMixin href "" "omnipy.compute._job.JobTemplateMixin"
              click omnipy.compute.flow.FlowBase href "" "omnipy.compute.flow.FlowBase"
            

Implement the core template behavior for function flows.

A function flow template wraps a Python callable that orchestrates work as a flow. Use this when the control flow is easiest to express directly in Python instead of as an explicit task list or dependency graph.

Decorator usage

Apply the template factory as a decorator to a Python callable. The wrapped callable becomes a reusable job template whose public outer signature is visible to template users and to the applied jobs created from it.

Wrapped callable

The wrapped callable defines both the implementation and the public outer signature of the flow.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.FuncFlowTemplate()
... def append_suffix_to_all(
...     dataset: TextDataset,
...     suffix: str,
... ) -> TextDataset:
...     output_dataset = TextDataset()
...     for title, data_file in dataset.items():
...         output_dataset[title] = f'{data_file.content}{suffix}'
...     return output_dataset
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> append_suffix_to_all.run(text_files, '!') == expected
True
Outer signature and modifiers

The wrapped callable's parameter list and return annotation define the outer interface of the template.

fixed_params permanently supplies selected callable parameters.

param_key_map renames selected callable parameters to external keyword names that callers or parent flows use when supplying inputs.

iterate_over_data_files, output_dataset_param, and output_dataset_cls adapt that outer interface for dataset-wise iteration.

When iterate_over_data_files=True and the inner first parameter is annotated as Model[T], callers see an outer dataset: Dataset[Model[T]] parameter and the outer return type becomes a dataset of the per-item return type. The inner callable still receives one model object at a time.

result_key wraps the returned value in a single-key dictionary, which is especially useful when a downstream DAG step should receive the result under a predictable name.

Examples:

>>> # With modifiers
>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_other(number: int, other: int) -> int:
...     return number + other
>>> plus_one = plus_other.refine(fixed_params={'other': 1})
>>> plus_one.run(4)
5
>>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
>>> plus_x.run(4, x=3)
7
>>> plus_one_dict = plus_one.refine(result_key='number')
>>> plus_one_dict.run(4)
{'number': 5}

Examples:

>>> # With dataset-wise iteration
>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def add_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> add_suffix.run(text_files, suffix='!') == expected
True
Tasks and flows

Tasks are terminal jobs: they wrap one callable and execute one compute step.

Flows are orchestration jobs: they may contain child tasks and child flows, so larger pipelines can be assembled hierarchically from smaller reusable pieces.

Lifecycle

Apply a template with apply() to create a runnable job with engine decorators and current config attached. Call the resulting applied job with runtime arguments.

Use run() as a shorthand for apply() followed immediately by calling the applied job.

Use refine() to reuse a template while changing configuration such as name, fixed_params, or param_key_map.

Use revise() on an applied job to reconstruct a template from that job's current configuration.

Examples:

>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_one(number: int) -> int:
...     return number + 1
>>> plus_one.run(1)
2
>>> applied_job = plus_one.apply()
>>> applied_job(2)
3
>>> refined_template = plus_one.refine(name='plus_one_renamed')
>>> revised_template = applied_job.revise()

Instances are normally produced through the FuncFlowTemplate decorator factory rather than by direct construction.

METHOD DESCRIPTION
__init__
accept_mixin

Register a mixin class for dynamic composition.

apply

Create an applied job from this template without executing it.

create_job_template

Create a job template instance from the concrete template class.

log

Emit a log message, optionally using an explicit event timestamp.

refine

Forward refinement to the shared template lifecycle implementation.

reset_mixins

Clear all accepted mixins and restore the original init signature.

run

Apply the template and execute the resulting job immediately.

ATTRIBUTE DESCRIPTION
callable_type

TYPE: CallableType.Literals

config

Return the job configuration visible to this instance.

TYPE: IsJobConfig

engine

Return the engine associated with this job, if any.

TYPE: IsEngine | None

in_flow_context

Return whether the job is currently executing inside a flow context.

TYPE: bool

logger

Return the logger bound to the concrete instance type.

TYPE: Logger

Source code in src/omnipy/compute/flow.py
class FuncFlowTemplateCore(
        FuncArgJobBase[
            IsFuncFlowTemplate[_CallP, _RetT],
            IsFuncFlow[_CallP, _RetT],
            _CallP,
            _RetT,
        ],
        JobTemplateMixin[
            IsFuncFlowTemplate[_CallP, _RetT],
            IsFuncFlow[_CallP, _RetT],
            _CallP,
            _RetT,
        ],
        FlowBase,
        Generic[_CallP, _RetT],
):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # Implement the core template behavior for function flows.
    #
    # {{FUNC_FLOW_TEMPLATE_DESCRIPTION}}
    #
    # Instances are normally produced through the [FuncFlowTemplate][] decorator
    # factory rather than by direct construction.
    #
    """Implement the core template behavior for function flows.

    A function flow template wraps a Python callable that orchestrates
    work as a flow. Use this when the control flow is easiest to
    express directly in Python instead of as an explicit task list or
    dependency graph.

    ### Decorator usage

    Apply the template factory as a decorator to a Python callable.
    The wrapped callable becomes a reusable job template whose public outer
    signature is visible to template users and to the applied jobs created
    from it.

    ### Wrapped callable

    The wrapped callable defines both the implementation and the public
    outer signature of the flow.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.FuncFlowTemplate()
        ... def append_suffix_to_all(
        ...     dataset: TextDataset,
        ...     suffix: str,
        ... ) -> TextDataset:
        ...     output_dataset = TextDataset()
        ...     for title, data_file in dataset.items():
        ...         output_dataset[title] = f'{data_file.content}{suffix}'
        ...     return output_dataset

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> append_suffix_to_all.run(text_files, '!') == expected
        True

    ### Outer signature and modifiers

    The wrapped callable's parameter list and return annotation define the
    outer interface of the template.

    ``fixed_params`` permanently supplies selected callable parameters.

    ``param_key_map`` renames selected callable parameters to external
    keyword names that callers or parent flows use when supplying inputs.

    ``iterate_over_data_files``, ``output_dataset_param``, and
    ``output_dataset_cls`` adapt that outer interface for dataset-wise
    iteration.

    When ``iterate_over_data_files=True`` and the inner first parameter is
    annotated as ``Model[T]``, callers see an outer
    ``dataset: Dataset[Model[T]]`` parameter and the outer return type
    becomes a dataset of the per-item return type. The inner callable still
    receives one model object at a time.

    ``result_key`` wraps the returned value in a single-key dictionary,
    which is especially useful when a downstream DAG step should receive
    the result under a predictable name.

    Examples:
        >>> # With modifiers
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_other(number: int, other: int) -> int:
        ...     return number + other

        >>> plus_one = plus_other.refine(fixed_params={'other': 1})
        >>> plus_one.run(4)
        5
        >>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
        >>> plus_x.run(4, x=3)
        7
        >>> plus_one_dict = plus_one.refine(result_key='number')
        >>> plus_one_dict.run(4)
        {'number': 5}

    Examples:
        >>> # With dataset-wise iteration
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> add_suffix.run(text_files, suffix='!') == expected
        True

    ### Tasks and flows

    Tasks are terminal jobs: they wrap one callable and execute one compute
    step.

    Flows are orchestration jobs: they may contain child tasks and child
    flows, so larger pipelines can be assembled hierarchically from smaller
    reusable pieces.

    ### Lifecycle

    Apply a template with [`apply()`][omnipy.compute._job.JobTemplateMixin.apply]
    to create a runnable job with engine decorators and current config attached.
    Call the resulting applied job with runtime arguments.

    Use [`run()`][omnipy.compute._job.JobTemplateMixin.run] as a shorthand for
    ``apply()`` followed immediately by calling the applied job.

    Use [`refine()`][omnipy.compute._job.JobTemplateMixin.refine] to reuse a
    template while changing configuration such as ``name``, ``fixed_params``,
    or ``param_key_map``.

    Use [`revise()`][omnipy.compute._job.JobMixin.revise] on an applied job to
    reconstruct a template from that job's current configuration.

    Examples:
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_one(number: int) -> int:
        ...     return number + 1

        >>> plus_one.run(1)
        2
        >>> applied_job = plus_one.apply()
        >>> applied_job(2)
        3
        >>> refined_template = plus_one.refine(name='plus_one_renamed')
        >>> revised_template = applied_job.revise()

    Instances are normally produced through the [FuncFlowTemplate][] decorator
    factory rather than by direct construction.
    """
    @classmethod
    def _get_job_subcls_for_apply(cls) -> type[IsFuncFlow[_CallP, _RetT]]:
        """Return the executable function flow type produced by this template.

        The template/application machinery calls this hook when it needs the
        concrete flow class that should be instantiated from a function flow
        template.

        Returns:
            type[IsFuncFlow[_CallP, _RetT]]: The executable [FuncFlow][]
                subclass associated with this template.
        """
        return cast(type[IsFuncFlow[_CallP, _RetT]], FuncFlow[_CallP, _RetT])

callable_type property

callable_type: CallableType.Literals

config property

config: IsJobConfig

Return the job configuration visible to this instance.

RETURNS DESCRIPTION
IsJobConfig

Active job configuration used for runtime behavior.

TYPE: IsJobConfig

engine property

engine: IsEngine | None

Return the engine associated with this job, if any.

RETURNS DESCRIPTION
IsEngine | None

IsEngine | None: Engine used for decoration and execution, or None.

in_flow_context property

in_flow_context: bool

Return whether the job is currently executing inside a flow context.

RETURNS DESCRIPTION
bool

True when a surrounding flow context is active.

TYPE: bool

logger property

logger: Logger

Return the logger bound to the concrete instance type.

RETURNS DESCRIPTION
Logger

Logger used by the object for Omnipy log messages.

TYPE: Logger

__init__

__init__(job_func: Callable[_CallP, _RetT], /, *args: object, **kwargs: object) -> None
Source code in src/omnipy/compute/_func_job.py
def __init__(self, job_func: Callable[_CallP, _RetT], /, *args: object,
             **kwargs: object) -> None:
    self._job_func = job_func

accept_mixin classmethod

accept_mixin(mixin_cls: Type) -> None

Register a mixin class for dynamic composition.

PARAMETER DESCRIPTION
mixin_cls

Mixin class whose __init__ keyword-only parameters should be merged into the acceptor signature.

TYPE: Type

Source code in src/omnipy/util/mixin.py
@classmethod
def accept_mixin(cls, mixin_cls: Type) -> None:
    """Register a mixin class for dynamic composition.

    Args:
        mixin_cls: Mixin class whose ``__init__`` keyword-only parameters
            should be merged into the acceptor signature.
    """
    cls._accept_mixin(mixin_cls, update=True)

apply

apply() -> _JobT

Create an applied job from this template without executing it.

RETURNS DESCRIPTION
_JobT

Applied job instance ready to be called.

TYPE: _JobT

Source code in src/omnipy/compute/_job.py
def apply(self) -> _JobT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_APPLY_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_APPLY_DETAILS}}
    """Create an applied job from this template without executing it.

    Returns:
        _JobT: Applied job instance ready to be called.
    """
    job = self._cast_to_job_tmpl()._apply()
    update_wrapper(job, self, updated=[])
    return cast(_JobT, job)

create_job_template classmethod

create_job_template(*args: object, **kwargs: object) -> _JobTemplateT

Create a job template instance from the concrete template class.

PARAMETER DESCRIPTION
*args

Positional constructor arguments.

TYPE: object DEFAULT: ()

**kwargs

Keyword constructor arguments.

TYPE: object DEFAULT: {}

RETURNS DESCRIPTION
_JobTemplateT

New job template instance.

TYPE: _JobTemplateT

Source code in src/omnipy/compute/_job.py
@classmethod
def create_job_template(cls, *args: object, **kwargs: object) -> _JobTemplateT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_CREATE_JOB_TEMPLATE_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_CREATE_JOB_TEMPLATE_DETAILS}}
    """Create a job template instance from the concrete template class.

    Args:
        *args: Positional constructor arguments.
        **kwargs: Keyword constructor arguments.

    Returns:
        _JobTemplateT: New job template instance.
    """

    cls_as_job_base = cast(type[IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT]], cls)
    return cls_as_job_base._create_job_template(*args, **kwargs)

log

log(log_msg: str, level: int = INFO, datetime_obj: datetime | None = None)

Emit a log message, optionally using an explicit event timestamp.

PARAMETER DESCRIPTION
log_msg

Message text to send to the logger.

TYPE: str

level

Standard library logging level.

TYPE: int DEFAULT: INFO

datetime_obj

Timestamp to attach to the record instead of wall-clock time.

TYPE: datetime | None DEFAULT: None

Source code in src/omnipy/hub/log/mixin.py
def log(self, log_msg: str, level: int = INFO, datetime_obj: datetime | None = None):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{CANLOG_LOG_SUMMARY}}
    #
    # {{CANLOG_LOG_DETAILS}}
    """Emit a log message, optionally using an explicit event timestamp.

    Args:
        log_msg: Message text to send to the logger.
        level: Standard library logging level.
        datetime_obj: Timestamp to attach to the record instead of wall-clock time.
    """
    if self._logger is not None:
        create_time = datetime_obj.timestamp() if datetime_obj else time.time()
        self._logger.log(level, log_msg, extra=dict(timestamp=create_time))

refine

refine(*args: Any, update: bool = True, **kwargs: object) -> _JobTemplateT

Forward refinement to the shared template lifecycle implementation.

See IsFuncArgJobTemplate.refine and IsChildJobListArgJobTemplate.refine.

Source code in src/omnipy/compute/_job.py
def refine(self, *args: Any, update: bool = True, **kwargs: object) -> _JobTemplateT:
    """Forward refinement to the shared template lifecycle implementation.

    See [`IsFuncArgJobTemplate.refine`]
    [omnipy.shared.protocols.compute.job.IsFuncArgJobTemplate.refine] and
    [`IsChildJobListArgJobTemplate.refine`]
    [omnipy.shared.protocols.compute.job.IsChildJobListArgJobTemplate.refine].
    """
    self_as_job_base = cast(IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT], self)
    return self_as_job_base._refine(*args, update=update, **kwargs)

reset_mixins classmethod

reset_mixins()

Clear all accepted mixins and restore the original init signature.

Source code in src/omnipy/util/mixin.py
@classmethod
def reset_mixins(cls):
    """Clear all accepted mixins and restore the original init signature."""
    cls._mixin_classes.clear()
    cls._init_params_per_mixin_cls.clear()
    cls.__init__.__signature__ = cls._orig_init_signature

run

run(*args: _CallP.args, **kwargs: _CallP.kwargs) -> _RetT

Apply the template and execute the resulting job immediately.

PARAMETER DESCRIPTION
*args

Positional arguments passed to the applied job.

TYPE: _CallP.args DEFAULT: ()

**kwargs

Keyword arguments passed to the applied job.

TYPE: _CallP.kwargs DEFAULT: {}

RETURNS DESCRIPTION
_RetCovT

Result returned by the applied job.

TYPE: _RetT

Source code in src/omnipy/compute/_job.py
def run(self, *args: _CallP.args, **kwargs: _CallP.kwargs) -> _RetT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_RUN_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_RUN_DETAILS}}
    """Apply the template and execute the resulting job immediately.

    Args:
        *args: Positional arguments passed to the applied job.
        **kwargs: Keyword arguments passed to the applied job.

    Returns:
        _RetCovT: Result returned by the applied job.
    """
    # TODO: Using JobTemplateMixin.run() inside flows should give error message
    return self._cast_to_job_tmpl().apply()(*args, **kwargs)

LinearFlow

Bases: JobMixin[IsLinearFlowTemplate[_CallP, _RetT], IsLinearFlow[_CallP, _RetT], _CallP, _RetT], FlowBase, ChildJobListArgJobBase[IsLinearFlowTemplate[_CallP, _RetT], IsLinearFlow[_CallP, _RetT], _CallP, _RetT], Generic[_CallP, _RetT]


              flowchart BT
              omnipy.compute.flow.LinearFlow[LinearFlow]
              omnipy.compute._job.JobMixin[JobMixin]
              omnipy.compute.flow.FlowBase[FlowBase]
              omnipy.compute._joblist_job.ChildJobListArgJobBase[ChildJobListArgJobBase]
              omnipy.compute._func_job.FuncArgJobBase[FuncArgJobBase]
              omnipy.compute._func_job.PlainFuncArgJobBase[PlainFuncArgJobBase]
              omnipy.compute._job.JobBase[JobBase]
              omnipy.hub.log.mixin.LogMixin[LogMixin]
              omnipy.util.mixin.DynamicMixinAcceptor[DynamicMixinAcceptor]

                              omnipy.compute._job.JobMixin --> omnipy.compute.flow.LinearFlow
                                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobMixin
                

                omnipy.compute.flow.FlowBase --> omnipy.compute.flow.LinearFlow
                
                omnipy.compute._joblist_job.ChildJobListArgJobBase --> omnipy.compute.flow.LinearFlow
                                omnipy.compute._func_job.FuncArgJobBase --> omnipy.compute._joblist_job.ChildJobListArgJobBase
                                omnipy.compute._func_job.PlainFuncArgJobBase --> omnipy.compute._func_job.FuncArgJobBase
                                omnipy.compute._job.JobBase --> omnipy.compute._func_job.PlainFuncArgJobBase
                                omnipy.hub.log.mixin.LogMixin --> omnipy.compute._job.JobBase
                
                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobBase
                






              click omnipy.compute.flow.LinearFlow href "" "omnipy.compute.flow.LinearFlow"
              click omnipy.compute._job.JobMixin href "" "omnipy.compute._job.JobMixin"
              click omnipy.compute.flow.FlowBase href "" "omnipy.compute.flow.FlowBase"
              click omnipy.compute._joblist_job.ChildJobListArgJobBase href "" "omnipy.compute._joblist_job.ChildJobListArgJobBase"
              click omnipy.compute._func_job.FuncArgJobBase href "" "omnipy.compute._func_job.FuncArgJobBase"
              click omnipy.compute._func_job.PlainFuncArgJobBase href "" "omnipy.compute._func_job.PlainFuncArgJobBase"
              click omnipy.compute._job.JobBase href "" "omnipy.compute._job.JobBase"
              click omnipy.hub.log.mixin.LogMixin href "" "omnipy.hub.log.mixin.LogMixin"
              click omnipy.util.mixin.DynamicMixinAcceptor href "" "omnipy.util.mixin.DynamicMixinAcceptor"
            

Execute a flow whose tasks run in declaration order.

A LinearFlow runs its constituent tasks in declaration order, making each step wait for the previous one to finish. Use it for pipelines where every stage depends on the output or side effects of the stage before it.

Instances are typically produced by calling a LinearFlowTemplate rather than by constructing LinearFlow directly.

CLASS DESCRIPTION
DataClassAndJobParentInfo
METHOD DESCRIPTION
__init__
accept_mixin

Register a mixin class for dynamic composition.

create_job

Create an applied job instance from the concrete job class.

log

Emit a log message, optionally using an explicit event timestamp.

reset_mixins

Clear all accepted mixins and restore the original init signature.

revise

Return a template reconstructed from this applied job.

ATTRIBUTE DESCRIPTION
callable_type

TYPE: CallableType.Literals

child_job_templates

TYPE: tuple[ChildJobTemplateLike, ...]

config

Return the job configuration visible to this instance.

TYPE: IsJobConfig

engine

Return the engine associated with this job, if any.

TYPE: IsEngine | None

in_flow_context

Return whether the job is currently executing inside a flow context.

TYPE: bool

logger

Return the logger bound to the concrete instance type.

TYPE: Logger

time_of_cur_toplevel_flow_run

Return the start time of the active top-level flow run, if any.

TYPE: datetime | None

Source code in src/omnipy/compute/flow.py
class LinearFlow(JobMixin[IsLinearFlowTemplate[_CallP, _RetT],
                          IsLinearFlow[_CallP, _RetT],
                          _CallP,
                          _RetT],
                 FlowBase,
                 ChildJobListArgJobBase[IsLinearFlowTemplate[_CallP, _RetT],
                                        IsLinearFlow[_CallP, _RetT],
                                        _CallP,
                                        _RetT],
                 Generic[_CallP, _RetT]):
    """Execute a flow whose tasks run in declaration order.

    A ``LinearFlow`` runs its constituent tasks in declaration order, making
    each step wait for the previous one to finish. Use it for pipelines where
    every stage depends on the output or side effects of the stage before it.

    Instances are typically produced by calling a ``LinearFlowTemplate``
    rather than by constructing ``LinearFlow`` directly.
    """
    def _apply_engine_decorator(self, engine: IsEngine) -> None:
        """Register the engine decorator for linear-flow execution.

        When the flow already holds a runner engine, this method asks that
        engine to wrap the flow's call function with the engine-specific
        linear-flow execution behavior.

        Args:
            self: Current linear flow instance.
            engine: Engine candidate supplied during job setup.
        """
        if self.engine:
            engine = cast(IsJobRunnerEngine, self.engine)
            self_with_mixins = cast(IsLinearFlow, self)
            engine.apply_job_decorator(
                JobType.LINEAR_FLOW,
                self_with_mixins,
                self._accept_call_func_decorator,
            )

    @classmethod
    def _get_job_template_subcls_for_revise(cls) -> type[IsLinearFlowTemplate[_CallP, _RetT]]:
        """Return the template type used to revise this flow.

        Revision operations use this hook to recover the decorator-backed
        template class corresponding to an executable linear flow instance.

        Returns:
            type[IsLinearFlowTemplate[_CallP, _RetT]]: The [LinearFlowTemplate][]
                class associated with this flow.
        """
        return cast(type[IsLinearFlowTemplate[_CallP, _RetT]], LinearFlowTemplateCore)

callable_type property

callable_type: CallableType.Literals

child_job_templates property

child_job_templates: tuple[ChildJobTemplateLike, ...]

config property

config: IsJobConfig

Return the job configuration visible to this instance.

RETURNS DESCRIPTION
IsJobConfig

Active job configuration used for runtime behavior.

TYPE: IsJobConfig

engine property

engine: IsEngine | None

Return the engine associated with this job, if any.

RETURNS DESCRIPTION
IsEngine | None

IsEngine | None: Engine used for decoration and execution, or None.

in_flow_context property

in_flow_context: bool

Return whether the job is currently executing inside a flow context.

RETURNS DESCRIPTION
bool

True when a surrounding flow context is active.

TYPE: bool

logger property

logger: Logger

Return the logger bound to the concrete instance type.

RETURNS DESCRIPTION
Logger

Logger used by the object for Omnipy log messages.

TYPE: Logger

time_of_cur_toplevel_flow_run property

time_of_cur_toplevel_flow_run: datetime | None

Return the start time of the active top-level flow run, if any.

RETURNS DESCRIPTION
datetime | None

datetime | None: Timestamp for the current outermost flow run, or None.

DataClassAndJobParentInfo

Bases: NamedTuple


              flowchart BT
              omnipy.compute.flow.LinearFlow.DataClassAndJobParentInfo[DataClassAndJobParentInfo]

              

              click omnipy.compute.flow.LinearFlow.DataClassAndJobParentInfo href "" "omnipy.compute.flow.LinearFlow.DataClassAndJobParentInfo"
            
ATTRIBUTE DESCRIPTION
dataset_or_model

TYPE: Literal['dataset', 'model']

parent_coerces_from_kwargs

TYPE: bool

Source code in src/omnipy/compute/_joblist_job.py
class DataClassAndJobParentInfo(NamedTuple):
    dataset_or_model: Literal['dataset', 'model']
    parent_coerces_from_kwargs: bool

dataset_or_model instance-attribute

dataset_or_model: Literal['dataset', 'model']

parent_coerces_from_kwargs instance-attribute

parent_coerces_from_kwargs: bool

__init__

__init__(*args, **kwargs)
Source code in src/omnipy/compute/_job.py
def __init__(self, *args, **kwargs):
    if JobBase not in self.__class__.__mro__:
        raise TypeError('JobMixin is not meant to be instantiated outside the context '
                        'of a JobBase subclass.')

accept_mixin classmethod

accept_mixin(mixin_cls: Type) -> None

Register a mixin class for dynamic composition.

PARAMETER DESCRIPTION
mixin_cls

Mixin class whose __init__ keyword-only parameters should be merged into the acceptor signature.

TYPE: Type

Source code in src/omnipy/util/mixin.py
@classmethod
def accept_mixin(cls, mixin_cls: Type) -> None:
    """Register a mixin class for dynamic composition.

    Args:
        mixin_cls: Mixin class whose ``__init__`` keyword-only parameters
            should be merged into the acceptor signature.
    """
    cls._accept_mixin(mixin_cls, update=True)

create_job classmethod

create_job(*args: object, **kwargs: object) -> _JobT

Create an applied job instance from the concrete job class.

PARAMETER DESCRIPTION
*args

Positional constructor arguments.

TYPE: object DEFAULT: ()

**kwargs

Keyword constructor arguments.

TYPE: object DEFAULT: {}

RETURNS DESCRIPTION
_JobT

New applied job instance.

TYPE: _JobT

Source code in src/omnipy/compute/_job.py
@classmethod
def create_job(cls, *args: object, **kwargs: object) -> _JobT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOB_CREATE_JOB_SUMMARY}}
    #
    # {{ISJOB_CREATE_JOB_DETAILS}}
    """Create an applied job instance from the concrete job class.

    Args:
        *args: Positional constructor arguments.
        **kwargs: Keyword constructor arguments.

    Returns:
        _JobT: New applied job instance.
    """
    cls_as_job_base = cast(IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT], cls)
    return cls_as_job_base._create_job(*args, **kwargs)

log

log(log_msg: str, level: int = INFO, datetime_obj: datetime | None = None)

Emit a log message, optionally using an explicit event timestamp.

PARAMETER DESCRIPTION
log_msg

Message text to send to the logger.

TYPE: str

level

Standard library logging level.

TYPE: int DEFAULT: INFO

datetime_obj

Timestamp to attach to the record instead of wall-clock time.

TYPE: datetime | None DEFAULT: None

Source code in src/omnipy/hub/log/mixin.py
def log(self, log_msg: str, level: int = INFO, datetime_obj: datetime | None = None):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{CANLOG_LOG_SUMMARY}}
    #
    # {{CANLOG_LOG_DETAILS}}
    """Emit a log message, optionally using an explicit event timestamp.

    Args:
        log_msg: Message text to send to the logger.
        level: Standard library logging level.
        datetime_obj: Timestamp to attach to the record instead of wall-clock time.
    """
    if self._logger is not None:
        create_time = datetime_obj.timestamp() if datetime_obj else time.time()
        self._logger.log(level, log_msg, extra=dict(timestamp=create_time))

reset_mixins classmethod

reset_mixins()

Clear all accepted mixins and restore the original init signature.

Source code in src/omnipy/util/mixin.py
@classmethod
def reset_mixins(cls):
    """Clear all accepted mixins and restore the original init signature."""
    cls._mixin_classes.clear()
    cls._init_params_per_mixin_cls.clear()
    cls.__init__.__signature__ = cls._orig_init_signature

revise

revise() -> _JobTemplateT

Return a template reconstructed from this applied job.

RETURNS DESCRIPTION
_JobTemplateT

Template carrying the current job configuration.

TYPE: _JobTemplateT

Source code in src/omnipy/compute/_job.py
def revise(self) -> _JobTemplateT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOB_REVISE_SUMMARY}}
    #
    # {{ISJOB_REVISE_DETAILS}}
    """Return a template reconstructed from this applied job.

    Returns:
        _JobTemplateT: Template carrying the current job configuration.
    """
    self_as_job_base = cast(
        IsJobBase[IsJobTemplate[_JobTemplateT, _JobT, _CallP, _RetT], _JobT, _CallP, _RetT],
        self)
    job_template = self_as_job_base._revise()
    update_wrapper(job_template, self, updated=[])
    return cast(_JobTemplateT, job_template)

LinearFlowTemplateCore

Bases: ChildJobListArgJobBase[IsLinearFlowTemplate[_CallP, _RetT], IsLinearFlow[_CallP, _RetT], _CallP, _RetT], JobTemplateMixin[IsLinearFlowTemplate[_CallP, _RetT], IsLinearFlow[_CallP, _RetT], _CallP, _RetT], FlowBase, Generic[_CallP, _RetT]


              flowchart BT
              omnipy.compute.flow.LinearFlowTemplateCore[LinearFlowTemplateCore]
              omnipy.compute._joblist_job.ChildJobListArgJobBase[ChildJobListArgJobBase]
              omnipy.compute._func_job.FuncArgJobBase[FuncArgJobBase]
              omnipy.compute._func_job.PlainFuncArgJobBase[PlainFuncArgJobBase]
              omnipy.compute._job.JobBase[JobBase]
              omnipy.hub.log.mixin.LogMixin[LogMixin]
              omnipy.util.mixin.DynamicMixinAcceptor[DynamicMixinAcceptor]
              omnipy.compute._job.JobTemplateMixin[JobTemplateMixin]
              omnipy.compute.flow.FlowBase[FlowBase]

                              omnipy.compute._joblist_job.ChildJobListArgJobBase --> omnipy.compute.flow.LinearFlowTemplateCore
                                omnipy.compute._func_job.FuncArgJobBase --> omnipy.compute._joblist_job.ChildJobListArgJobBase
                                omnipy.compute._func_job.PlainFuncArgJobBase --> omnipy.compute._func_job.FuncArgJobBase
                                omnipy.compute._job.JobBase --> omnipy.compute._func_job.PlainFuncArgJobBase
                                omnipy.hub.log.mixin.LogMixin --> omnipy.compute._job.JobBase
                
                omnipy.util.mixin.DynamicMixinAcceptor --> omnipy.compute._job.JobBase
                




                omnipy.compute._job.JobTemplateMixin --> omnipy.compute.flow.LinearFlowTemplateCore
                
                omnipy.compute.flow.FlowBase --> omnipy.compute.flow.LinearFlowTemplateCore
                


              click omnipy.compute.flow.LinearFlowTemplateCore href "" "omnipy.compute.flow.LinearFlowTemplateCore"
              click omnipy.compute._joblist_job.ChildJobListArgJobBase href "" "omnipy.compute._joblist_job.ChildJobListArgJobBase"
              click omnipy.compute._func_job.FuncArgJobBase href "" "omnipy.compute._func_job.FuncArgJobBase"
              click omnipy.compute._func_job.PlainFuncArgJobBase href "" "omnipy.compute._func_job.PlainFuncArgJobBase"
              click omnipy.compute._job.JobBase href "" "omnipy.compute._job.JobBase"
              click omnipy.hub.log.mixin.LogMixin href "" "omnipy.hub.log.mixin.LogMixin"
              click omnipy.util.mixin.DynamicMixinAcceptor href "" "omnipy.util.mixin.DynamicMixinAcceptor"
              click omnipy.compute._job.JobTemplateMixin href "" "omnipy.compute._job.JobTemplateMixin"
              click omnipy.compute.flow.FlowBase href "" "omnipy.compute.flow.FlowBase"
            

Implement the core template behavior for linear flows.

A linear flow template wraps a Python callable together with an ordered list of child job templates that run sequentially. Use this when work should proceed step by step in a fixed declaration order.

Decorator usage

Apply the template factory as a decorator to a Python callable. The wrapped callable becomes a reusable job template whose public outer signature is visible to template users and to the applied jobs created from it.

Outer callable and child jobs

The wrapped callable defines the public outer signature of the flow, while the child-job list defines the executed body.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def normalize_text(data_file: TextModel) -> TextModel:
...     return data_file.content.strip().lower()
>>> @om.LinearFlowTemplate(normalize_text)
... def clean_texts(
...     dataset: TextDataset,
... ) -> TextDataset:
...     return dataset
>>> text_files = TextDataset({'a': ' Hi ', 'b': 'BYE '})
>>> expected = TextDataset({'a': 'hi', 'b': 'bye'})
>>> clean_texts.run(text_files) == expected
True
Child jobs and data classes

Child-job templates may be TaskTemplate or flow-template instances. This lets flows nest other flows as well as terminal tasks.

In linear and DAG flows, Model and Dataset subclasses are also allowed as child-job entries. During apply, they are coerced into helper task templates that construct the requested data object from positional input in linear flows or from keyword-matched input in DAG flows.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.LinearFlowTemplate(TextModel)
... def wrap_text(raw_text: str) -> TextModel:
...     return TextModel(raw_text)
>>> wrap_text.run('hello').content
'hello'
>>> @om.DagFlowTemplate(TextDataset)
... def collect_texts(first: str, second: str) -> TextDataset:
...     return TextDataset({'first': first, 'second': second})
>>> collect_texts.run(first='hello', second='bye') == TextDataset({
...     'first': 'hello',
...     'second': 'bye',
... })
True
Callable-type validation

For linear and DAG flows, the outer callable is primarily declarative: its signature exposes the public flow interface, while child jobs define the executed body.

The outer callable type is validated against the child-job composition. The terminal child determines whether the flow behaves like a function or generator, and any async child lifts the full flow to an async callable type.

As a result, a sync-function outer callable fits sync child execution, a generator outer callable fits generator-producing terminal children, and async outer callables are required when child composition is async. Mismatches raise TypeError when the flow template is created.

Generator shorthand with Void

When the validated outer callable must be a generator or async-generator only to expose the correct public signature, use Void in the body:

Examples:

>>> import omnipy as om
>>> from collections.abc import Iterator
>>> @om.TaskTemplate()
... def emit_lines() -> Iterator[str]:
...     yield 'first'
...     yield 'second'
>>> @om.LinearFlowTemplate(emit_lines)
... def line_stream() -> Iterator[str]:
...     yield from om.Void()

This shorthand exists only to satisfy the declared outer callable type; the child jobs still perform the actual flow work.

Outer signature and modifiers

The wrapped callable's parameter list and return annotation define the outer interface of the template.

fixed_params permanently supplies selected callable parameters.

param_key_map renames selected callable parameters to external keyword names that callers or parent flows use when supplying inputs.

iterate_over_data_files, output_dataset_param, and output_dataset_cls adapt that outer interface for dataset-wise iteration.

When iterate_over_data_files=True and the inner first parameter is annotated as Model[T], callers see an outer dataset: Dataset[Model[T]] parameter and the outer return type becomes a dataset of the per-item return type. The inner callable still receives one model object at a time.

result_key wraps the returned value in a single-key dictionary, which is especially useful when a downstream DAG step should receive the result under a predictable name.

Examples:

>>> # With modifiers
>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_other(number: int, other: int) -> int:
...     return number + other
>>> plus_one = plus_other.refine(fixed_params={'other': 1})
>>> plus_one.run(4)
5
>>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
>>> plus_x.run(4, x=3)
7
>>> plus_one_dict = plus_one.refine(result_key='number')
>>> plus_one_dict.run(4)
{'number': 5}

Examples:

>>> # With dataset-wise iteration
>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def add_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> add_suffix.run(text_files, suffix='!') == expected
True
Linear orchestration

Linear flows run child jobs strictly in declaration order.

The first child receives the caller's positional and keyword inputs. Each later child receives the previous child result as its leading positional input, plus any matching keyword arguments from the outer flow call.

fixed_params always override caller-supplied values for the child where they are configured. param_key_map lets later children expose different external keyword names than their underlying callable parameter names.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True)
... def strip_text(data_file: TextModel) -> TextModel:
...     return data_file.content.strip()
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def add_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> suffix_each = add_suffix.refine(param_key_map={'suffix': 'ending'})
>>> @om.LinearFlowTemplate(strip_text, suffix_each)
... def linear_flow(
...     dataset: TextDataset,
...     ending: str,
... ) -> TextDataset:
...     return dataset
>>> text_files = TextDataset({'a': ' hi', 'b': 'bye '})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> linear_flow.run(text_files, ending='!') == expected
True
Tasks and flows

Tasks are terminal jobs: they wrap one callable and execute one compute step.

Flows are orchestration jobs: they may contain child tasks and child flows, so larger pipelines can be assembled hierarchically from smaller reusable pieces.

Lifecycle

Apply a template with apply() to create a runnable job with engine decorators and current config attached. Call the resulting applied job with runtime arguments.

Use run() as a shorthand for apply() followed immediately by calling the applied job.

Use refine() to reuse a template while changing configuration such as name, fixed_params, or param_key_map.

Use revise() on an applied job to reconstruct a template from that job's current configuration.

Examples:

>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_one(number: int) -> int:
...     return number + 1
>>> plus_one.run(1)
2
>>> applied_job = plus_one.apply()
>>> applied_job(2)
3
>>> refined_template = plus_one.refine(name='plus_one_renamed')
>>> revised_template = applied_job.revise()

Instances are normally produced through the LinearFlowTemplate decorator factory rather than by direct construction.

CLASS DESCRIPTION
DataClassAndJobParentInfo
METHOD DESCRIPTION
__init__
accept_mixin

Register a mixin class for dynamic composition.

apply

Create an applied job from this template without executing it.

create_job_template

Create a job template instance from the concrete template class.

log

Emit a log message, optionally using an explicit event timestamp.

refine

Forward refinement to the shared template lifecycle implementation.

reset_mixins

Clear all accepted mixins and restore the original init signature.

run

Apply the template and execute the resulting job immediately.

ATTRIBUTE DESCRIPTION
callable_type

TYPE: CallableType.Literals

child_job_templates

TYPE: tuple[ChildJobTemplateLike, ...]

config

Return the job configuration visible to this instance.

TYPE: IsJobConfig

engine

Return the engine associated with this job, if any.

TYPE: IsEngine | None

in_flow_context

Return whether the job is currently executing inside a flow context.

TYPE: bool

logger

Return the logger bound to the concrete instance type.

TYPE: Logger

Source code in src/omnipy/compute/flow.py
class LinearFlowTemplateCore(
        ChildJobListArgJobBase[
            IsLinearFlowTemplate[_CallP, _RetT],
            IsLinearFlow[_CallP, _RetT],
            _CallP,
            _RetT,
        ],
        JobTemplateMixin[
            IsLinearFlowTemplate[_CallP, _RetT],
            IsLinearFlow[_CallP, _RetT],
            _CallP,
            _RetT,
        ],
        FlowBase,
        Generic[_CallP, _RetT],
):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # Implement the core template behavior for linear flows.
    #
    # {{LINEAR_FLOW_TEMPLATE_DESCRIPTION}}
    #
    # Instances are normally produced through the [LinearFlowTemplate][]
    # decorator factory rather than by direct construction.
    #
    """Implement the core template behavior for linear flows.

    A linear flow template wraps a Python callable together with an ordered
    list of child job templates that run sequentially. Use this when work
    should proceed step by step in a fixed declaration order.

    ### Decorator usage

    Apply the template factory as a decorator to a Python callable.
    The wrapped callable becomes a reusable job template whose public outer
    signature is visible to template users and to the applied jobs created
    from it.

    ### Outer callable and child jobs

    The wrapped callable defines the public outer signature of the flow,
    while the child-job list defines the executed body.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def normalize_text(data_file: TextModel) -> TextModel:
        ...     return data_file.content.strip().lower()

        >>> @om.LinearFlowTemplate(normalize_text)
        ... def clean_texts(
        ...     dataset: TextDataset,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': ' Hi ', 'b': 'BYE '})
        >>> expected = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> clean_texts.run(text_files) == expected
        True

    ### Child jobs and data classes

    Child-job templates may be
    [TaskTemplate][omnipy.compute.task.TaskTemplate] or flow-template
    instances.
    This lets flows nest other flows as well as terminal tasks.

    In linear and DAG flows, Model and Dataset subclasses are also allowed
    as child-job entries. During apply, they are coerced into helper task
    templates that construct the requested data object from positional
    input in linear flows or from keyword-matched input in DAG flows.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.LinearFlowTemplate(TextModel)
        ... def wrap_text(raw_text: str) -> TextModel:
        ...     return TextModel(raw_text)

        >>> wrap_text.run('hello').content
        'hello'
        >>> @om.DagFlowTemplate(TextDataset)
        ... def collect_texts(first: str, second: str) -> TextDataset:
        ...     return TextDataset({'first': first, 'second': second})

        >>> collect_texts.run(first='hello', second='bye') == TextDataset({
        ...     'first': 'hello',
        ...     'second': 'bye',
        ... })
        True

    ### Callable-type validation

    For linear and DAG flows, the outer callable is primarily declarative:
    its signature exposes the public flow interface, while child jobs define
    the executed body.

    The outer callable type is validated against the child-job composition.
    The terminal child determines whether the flow behaves like a function
    or generator, and any async child lifts the full flow to an async
    callable type.

    As a result, a sync-function outer callable fits sync child execution,
    a generator outer callable fits generator-producing terminal children,
    and async outer callables are required when child composition is async.
    Mismatches raise ``TypeError`` when the flow template is created.

    ### Generator shorthand with ``Void``

    When the validated outer callable must be a generator or
    async-generator only to expose the correct public signature, use
    [Void][omnipy.compute.helpers.Void] in the body:

    Examples:
        >>> import omnipy as om
        >>> from collections.abc import Iterator

        >>> @om.TaskTemplate()
        ... def emit_lines() -> Iterator[str]:
        ...     yield 'first'
        ...     yield 'second'

        >>> @om.LinearFlowTemplate(emit_lines)
        ... def line_stream() -> Iterator[str]:
        ...     yield from om.Void()

    This shorthand exists only to satisfy the declared outer callable type;
    the child jobs still perform the actual flow work.

    ### Outer signature and modifiers

    The wrapped callable's parameter list and return annotation define the
    outer interface of the template.

    ``fixed_params`` permanently supplies selected callable parameters.

    ``param_key_map`` renames selected callable parameters to external
    keyword names that callers or parent flows use when supplying inputs.

    ``iterate_over_data_files``, ``output_dataset_param``, and
    ``output_dataset_cls`` adapt that outer interface for dataset-wise
    iteration.

    When ``iterate_over_data_files=True`` and the inner first parameter is
    annotated as ``Model[T]``, callers see an outer
    ``dataset: Dataset[Model[T]]`` parameter and the outer return type
    becomes a dataset of the per-item return type. The inner callable still
    receives one model object at a time.

    ``result_key`` wraps the returned value in a single-key dictionary,
    which is especially useful when a downstream DAG step should receive
    the result under a predictable name.

    Examples:
        >>> # With modifiers
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_other(number: int, other: int) -> int:
        ...     return number + other

        >>> plus_one = plus_other.refine(fixed_params={'other': 1})
        >>> plus_one.run(4)
        5
        >>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
        >>> plus_x.run(4, x=3)
        7
        >>> plus_one_dict = plus_one.refine(result_key='number')
        >>> plus_one_dict.run(4)
        {'number': 5}

    Examples:
        >>> # With dataset-wise iteration
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> add_suffix.run(text_files, suffix='!') == expected
        True

    ### Linear orchestration

    Linear flows run child jobs strictly in declaration order.

    The first child receives the caller's positional and keyword inputs.
    Each later child receives the previous child result as its leading
    positional input, plus any matching keyword arguments from the outer
    flow call.

    ``fixed_params`` always override caller-supplied values for the child
    where they are configured. ``param_key_map`` lets later children expose
    different external keyword names than their underlying callable
    parameter names.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True)
        ... def strip_text(data_file: TextModel) -> TextModel:
        ...     return data_file.content.strip()

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> suffix_each = add_suffix.refine(param_key_map={'suffix': 'ending'})
        >>> @om.LinearFlowTemplate(strip_text, suffix_each)
        ... def linear_flow(
        ...     dataset: TextDataset,
        ...     ending: str,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': ' hi', 'b': 'bye '})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> linear_flow.run(text_files, ending='!') == expected
        True

    ### Tasks and flows

    Tasks are terminal jobs: they wrap one callable and execute one compute
    step.

    Flows are orchestration jobs: they may contain child tasks and child
    flows, so larger pipelines can be assembled hierarchically from smaller
    reusable pieces.

    ### Lifecycle

    Apply a template with [`apply()`][omnipy.compute._job.JobTemplateMixin.apply]
    to create a runnable job with engine decorators and current config attached.
    Call the resulting applied job with runtime arguments.

    Use [`run()`][omnipy.compute._job.JobTemplateMixin.run] as a shorthand for
    ``apply()`` followed immediately by calling the applied job.

    Use [`refine()`][omnipy.compute._job.JobTemplateMixin.refine] to reuse a
    template while changing configuration such as ``name``, ``fixed_params``,
    or ``param_key_map``.

    Use [`revise()`][omnipy.compute._job.JobMixin.revise] on an applied job to
    reconstruct a template from that job's current configuration.

    Examples:
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_one(number: int) -> int:
        ...     return number + 1

        >>> plus_one.run(1)
        2
        >>> applied_job = plus_one.apply()
        >>> applied_job(2)
        3
        >>> refined_template = plus_one.refine(name='plus_one_renamed')
        >>> revised_template = applied_job.revise()

    Instances are normally produced through the [LinearFlowTemplate][]
    decorator factory rather than by direct construction.
    """
    @classmethod
    def _get_job_subcls_for_apply(cls) -> type[IsLinearFlow[_CallP, _RetT]]:
        """Return the executable flow type produced by this template.

        The template/application machinery calls this hook when it needs the
        concrete job class that should be instantiated from a linear flow
        template.

        Returns:
            type[IsLinearFlow[_CallP, _RetT]]: The executable [LinearFlow][]
                subclass associated with this template.
        """
        return cast(type[IsLinearFlow[_CallP, _RetT]], LinearFlow[_CallP, _RetT])

callable_type property

callable_type: CallableType.Literals

child_job_templates property

child_job_templates: tuple[ChildJobTemplateLike, ...]

config property

config: IsJobConfig

Return the job configuration visible to this instance.

RETURNS DESCRIPTION
IsJobConfig

Active job configuration used for runtime behavior.

TYPE: IsJobConfig

engine property

engine: IsEngine | None

Return the engine associated with this job, if any.

RETURNS DESCRIPTION
IsEngine | None

IsEngine | None: Engine used for decoration and execution, or None.

in_flow_context property

in_flow_context: bool

Return whether the job is currently executing inside a flow context.

RETURNS DESCRIPTION
bool

True when a surrounding flow context is active.

TYPE: bool

logger property

logger: Logger

Return the logger bound to the concrete instance type.

RETURNS DESCRIPTION
Logger

Logger used by the object for Omnipy log messages.

TYPE: Logger

DataClassAndJobParentInfo

Bases: NamedTuple


              flowchart BT
              omnipy.compute.flow.LinearFlowTemplateCore.DataClassAndJobParentInfo[DataClassAndJobParentInfo]

              

              click omnipy.compute.flow.LinearFlowTemplateCore.DataClassAndJobParentInfo href "" "omnipy.compute.flow.LinearFlowTemplateCore.DataClassAndJobParentInfo"
            
ATTRIBUTE DESCRIPTION
dataset_or_model

TYPE: Literal['dataset', 'model']

parent_coerces_from_kwargs

TYPE: bool

Source code in src/omnipy/compute/_joblist_job.py
class DataClassAndJobParentInfo(NamedTuple):
    dataset_or_model: Literal['dataset', 'model']
    parent_coerces_from_kwargs: bool

dataset_or_model instance-attribute

dataset_or_model: Literal['dataset', 'model']

parent_coerces_from_kwargs instance-attribute

parent_coerces_from_kwargs: bool

__init__

__init__(
    job_func: Callable[_CallP, _RetT],
    /,
    *child_job_templates: ChildJobTemplateLike,
    **kwargs: object,
) -> None
Source code in src/omnipy/compute/_joblist_job.py
def __init__(self,
             job_func: Callable[_CallP, _RetT],
             /,
             *child_job_templates: ChildJobTemplateLike,
             **kwargs: object) -> None:
    super().__init__(job_func, *child_job_templates, **kwargs)
    self._child_job_templates: tuple[ChildJobTemplateLike, ...] = child_job_templates
    self._validate_callable_type_against_child_job_templates()

accept_mixin classmethod

accept_mixin(mixin_cls: Type) -> None

Register a mixin class for dynamic composition.

PARAMETER DESCRIPTION
mixin_cls

Mixin class whose __init__ keyword-only parameters should be merged into the acceptor signature.

TYPE: Type

Source code in src/omnipy/util/mixin.py
@classmethod
def accept_mixin(cls, mixin_cls: Type) -> None:
    """Register a mixin class for dynamic composition.

    Args:
        mixin_cls: Mixin class whose ``__init__`` keyword-only parameters
            should be merged into the acceptor signature.
    """
    cls._accept_mixin(mixin_cls, update=True)

apply

apply() -> _JobT

Create an applied job from this template without executing it.

RETURNS DESCRIPTION
_JobT

Applied job instance ready to be called.

TYPE: _JobT

Source code in src/omnipy/compute/_job.py
def apply(self) -> _JobT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_APPLY_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_APPLY_DETAILS}}
    """Create an applied job from this template without executing it.

    Returns:
        _JobT: Applied job instance ready to be called.
    """
    job = self._cast_to_job_tmpl()._apply()
    update_wrapper(job, self, updated=[])
    return cast(_JobT, job)

create_job_template classmethod

create_job_template(*args: object, **kwargs: object) -> _JobTemplateT

Create a job template instance from the concrete template class.

PARAMETER DESCRIPTION
*args

Positional constructor arguments.

TYPE: object DEFAULT: ()

**kwargs

Keyword constructor arguments.

TYPE: object DEFAULT: {}

RETURNS DESCRIPTION
_JobTemplateT

New job template instance.

TYPE: _JobTemplateT

Source code in src/omnipy/compute/_job.py
@classmethod
def create_job_template(cls, *args: object, **kwargs: object) -> _JobTemplateT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_CREATE_JOB_TEMPLATE_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_CREATE_JOB_TEMPLATE_DETAILS}}
    """Create a job template instance from the concrete template class.

    Args:
        *args: Positional constructor arguments.
        **kwargs: Keyword constructor arguments.

    Returns:
        _JobTemplateT: New job template instance.
    """

    cls_as_job_base = cast(type[IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT]], cls)
    return cls_as_job_base._create_job_template(*args, **kwargs)

log

log(log_msg: str, level: int = INFO, datetime_obj: datetime | None = None)

Emit a log message, optionally using an explicit event timestamp.

PARAMETER DESCRIPTION
log_msg

Message text to send to the logger.

TYPE: str

level

Standard library logging level.

TYPE: int DEFAULT: INFO

datetime_obj

Timestamp to attach to the record instead of wall-clock time.

TYPE: datetime | None DEFAULT: None

Source code in src/omnipy/hub/log/mixin.py
def log(self, log_msg: str, level: int = INFO, datetime_obj: datetime | None = None):
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{CANLOG_LOG_SUMMARY}}
    #
    # {{CANLOG_LOG_DETAILS}}
    """Emit a log message, optionally using an explicit event timestamp.

    Args:
        log_msg: Message text to send to the logger.
        level: Standard library logging level.
        datetime_obj: Timestamp to attach to the record instead of wall-clock time.
    """
    if self._logger is not None:
        create_time = datetime_obj.timestamp() if datetime_obj else time.time()
        self._logger.log(level, log_msg, extra=dict(timestamp=create_time))

refine

refine(*args: Any, update: bool = True, **kwargs: object) -> _JobTemplateT

Forward refinement to the shared template lifecycle implementation.

See IsFuncArgJobTemplate.refine and IsChildJobListArgJobTemplate.refine.

Source code in src/omnipy/compute/_job.py
def refine(self, *args: Any, update: bool = True, **kwargs: object) -> _JobTemplateT:
    """Forward refinement to the shared template lifecycle implementation.

    See [`IsFuncArgJobTemplate.refine`]
    [omnipy.shared.protocols.compute.job.IsFuncArgJobTemplate.refine] and
    [`IsChildJobListArgJobTemplate.refine`]
    [omnipy.shared.protocols.compute.job.IsChildJobListArgJobTemplate.refine].
    """
    self_as_job_base = cast(IsJobBase[_JobTemplateT, _JobT, _CallP, _RetT], self)
    return self_as_job_base._refine(*args, update=update, **kwargs)

reset_mixins classmethod

reset_mixins()

Clear all accepted mixins and restore the original init signature.

Source code in src/omnipy/util/mixin.py
@classmethod
def reset_mixins(cls):
    """Clear all accepted mixins and restore the original init signature."""
    cls._mixin_classes.clear()
    cls._init_params_per_mixin_cls.clear()
    cls.__init__.__signature__ = cls._orig_init_signature

run

run(*args: _CallP.args, **kwargs: _CallP.kwargs) -> _RetT

Apply the template and execute the resulting job immediately.

PARAMETER DESCRIPTION
*args

Positional arguments passed to the applied job.

TYPE: _CallP.args DEFAULT: ()

**kwargs

Keyword arguments passed to the applied job.

TYPE: _CallP.kwargs DEFAULT: {}

RETURNS DESCRIPTION
_RetCovT

Result returned by the applied job.

TYPE: _RetT

Source code in src/omnipy/compute/_job.py
def run(self, *args: _CallP.args, **kwargs: _CallP.kwargs) -> _RetT:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # {{ISJOBTEMPLATE_RUN_SUMMARY}}
    #
    # {{ISJOBTEMPLATE_RUN_DETAILS}}
    """Apply the template and execute the resulting job immediately.

    Args:
        *args: Positional arguments passed to the applied job.
        **kwargs: Keyword arguments passed to the applied job.

    Returns:
        _RetCovT: Result returned by the applied job.
    """
    # TODO: Using JobTemplateMixin.run() inside flows should give error message
    return self._cast_to_job_tmpl().apply()(*args, **kwargs)

DagFlowTemplate

DagFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    consume_kwargs_from_results: bool = True,
    iterate_over_data_files: Literal[True],
    output_dataset_cls: type[_RetDatasetClsT],
    **kwargs: Unpack[JobCommonKwargs],
) -> DagFlowTemplateIterWithDatasetClsDecorator[_RetDatasetClsT]
DagFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    consume_kwargs_from_results: bool = True,
    iterate_over_data_files: Literal[True],
    output_dataset_cls: None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> DagFlowTemplateIterDecorator
DagFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    consume_kwargs_from_results: bool = True,
    iterate_over_data_files: Literal[False] = False,
    output_dataset_cls: type[IsDataset] | None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> DagFlowTemplatePlainDecorator
DagFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    consume_kwargs_from_results: bool = True,
    iterate_over_data_files: bool,
    output_dataset_cls: type[_RetDatasetClsT] | None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> (
    DagFlowTemplatePlainDecorator
    | DagFlowTemplateIterDecorator
    | DagFlowTemplateIterWithDatasetClsDecorator[_RetDatasetClsT]
)

Decorator-style factory for defining directed acyclic graph flows.

A DAG flow template wraps a Python callable together with child job templates whose dependencies form a directed acyclic graph. Use this when flow steps branch and join but must not form cycles.

Decorator usage

Apply the template factory as a decorator to a Python callable. The wrapped callable becomes a reusable job template whose public outer signature is visible to template users and to the applied jobs created from it.

Outer callable and child jobs

The wrapped callable defines the public outer signature of the flow, while the child-job list defines the executed body.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True)
... def uppercase(data_file: TextModel) -> TextModel:
...     return data_file.content.upper()
>>> @om.TaskTemplate()
... def join_texts(
...     upper: TextDataset,
...     original: TextDataset,
... ) -> TextDataset:
...     merged = TextDataset()
...     for title in upper:
...         merged[title] = f'{upper[title].content}|{original[title].content}'
...     return merged
>>> @om.DagFlowTemplate(
...     uppercase.refine(result_key='upper'),
...     join_texts.refine(param_key_map={'upper': 'upper', 'original': 'dataset'}),
... )
... def my_dag(
...     dataset: TextDataset,
... ) -> TextDataset:
...     return dataset
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'HI|hi', 'b': 'BYE|bye'})
>>> my_dag.run(text_files) == expected
True
Child jobs and data classes

Child-job templates may be TaskTemplate or flow-template instances. This lets flows nest other flows as well as terminal tasks.

In linear and DAG flows, Model and Dataset subclasses are also allowed as child-job entries. During apply, they are coerced into helper task templates that construct the requested data object from positional input in linear flows or from keyword-matched input in DAG flows.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.LinearFlowTemplate(TextModel)
... def wrap_text(raw_text: str) -> TextModel:
...     return TextModel(raw_text)
>>> wrap_text.run('hello').content
'hello'
>>> @om.DagFlowTemplate(TextDataset)
... def collect_texts(first: str, second: str) -> TextDataset:
...     return TextDataset({'first': first, 'second': second})
>>> collect_texts.run(first='hello', second='bye') == TextDataset({
...     'first': 'hello',
...     'second': 'bye',
... })
True
Callable-type validation

For linear and DAG flows, the outer callable is primarily declarative: its signature exposes the public flow interface, while child jobs define the executed body.

The outer callable type is validated against the child-job composition. The terminal child determines whether the flow behaves like a function or generator, and any async child lifts the full flow to an async callable type.

As a result, a sync-function outer callable fits sync child execution, a generator outer callable fits generator-producing terminal children, and async outer callables are required when child composition is async. Mismatches raise TypeError when the flow template is created.

Generator shorthand with Void

When the validated outer callable must be a generator or async-generator only to expose the correct public signature, use Void in the body:

Examples:

>>> import omnipy as om
>>> from collections.abc import Iterator
>>> @om.TaskTemplate()
... def emit_lines() -> Iterator[str]:
...     yield 'first'
...     yield 'second'
>>> @om.LinearFlowTemplate(emit_lines)
... def line_stream() -> Iterator[str]:
...     yield from om.Void()

This shorthand exists only to satisfy the declared outer callable type; the child jobs still perform the actual flow work.

Outer signature and modifiers

The wrapped callable's parameter list and return annotation define the outer interface of the template.

fixed_params permanently supplies selected callable parameters.

param_key_map renames selected callable parameters to external keyword names that callers or parent flows use when supplying inputs.

iterate_over_data_files, output_dataset_param, and output_dataset_cls adapt that outer interface for dataset-wise iteration.

When iterate_over_data_files=True and the inner first parameter is annotated as Model[T], callers see an outer dataset: Dataset[Model[T]] parameter and the outer return type becomes a dataset of the per-item return type. The inner callable still receives one model object at a time.

result_key wraps the returned value in a single-key dictionary, which is especially useful when a downstream DAG step should receive the result under a predictable name.

Examples:

>>> # With modifiers
>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_other(number: int, other: int) -> int:
...     return number + other
>>> plus_one = plus_other.refine(fixed_params={'other': 1})
>>> plus_one.run(4)
5
>>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
>>> plus_x.run(4, x=3)
7
>>> plus_one_dict = plus_one.refine(result_key='number')
>>> plus_one_dict.run(4)
{'number': 5}

Examples:

>>> # With dataset-wise iteration
>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def add_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> add_suffix.run(text_files, suffix='!') == expected
True
DAG orchestration

DAG flows route values by keyword name instead of chaining every child result positionally.

The outer flow call first binds its arguments to the outer callable signature. Each child then receives the matching keyword subset from the accumulated named results.

By default, a non-dictionary child result is stored under the child job name. result_key stores it under a custom key instead, which is the usual way to make one branch feed another. Dictionary results merge directly into the accumulated named result set.

param_key_map lets a child read from externally visible DAG keys using different internal callable parameter names, and fixed_params pins selected child inputs regardless of what earlier branches produce.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate()
... def uppercase(data_file: TextModel) -> TextModel:
...     return data_file.content.upper()
>>> @om.TaskTemplate()
... def append_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> @om.TaskTemplate()
... def combine_texts(
...     left_dataset: TextDataset,
...     right_dataset: TextDataset,
... ) -> TextDataset:
...     merged = TextDataset()
...     for title in left_dataset:
...         merged[title] = (
...             f'{left_dataset[title].content}|'
...             f'{right_dataset[title].content}'
...         )
...     return merged
>>> @om.DagFlowTemplate(
...     uppercase.refine(
...         iterate_over_data_files=True,
...         result_key='upper',
...     ),
...     append_suffix.refine(
...         iterate_over_data_files=True,
...         result_key='suffixed',
...         param_key_map={'suffix': 'ending'},
...     ),
...     combine_texts.refine(
...         param_key_map={
...             'left_dataset': 'upper',
...             'right_dataset': 'suffixed',
...         },
...     ),
... )
... def dag_flow(
...     dataset: TextDataset,
...     ending: str,
... ) -> TextDataset:
...     return dataset
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({
...     'a': 'HI|hi!',
...     'b': 'BYE|bye!',
... })
>>> dag_flow.run(text_files, ending='!') == expected
True
Tasks and flows

Tasks are terminal jobs: they wrap one callable and execute one compute step.

Flows are orchestration jobs: they may contain child tasks and child flows, so larger pipelines can be assembled hierarchically from smaller reusable pieces.

Lifecycle

Apply a template with apply() to create a runnable job with engine decorators and current config attached. Call the resulting applied job with runtime arguments.

Use run() as a shorthand for apply() followed immediately by calling the applied job.

Use refine() to reuse a template while changing configuration such as name, fixed_params, or param_key_map.

Use revise() on an applied job to reconstruct a template from that job's current configuration.

Examples:

>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_one(number: int) -> int:
...     return number + 1
>>> plus_one.run(1)
2
>>> applied_job = plus_one.apply()
>>> applied_job(2)
3
>>> refined_template = plus_one.refine(name='plus_one_renamed')
>>> revised_template = applied_job.revise()
PARAMETER DESCRIPTION
*child_job_templates

Ordered templates of child jobs to be run as part of the parent job. Model and Dataset subclasses are also allowed, in which case they are converted to create_X_from_args (if linear flow) or create_X_from_kwargs (if DAG flow) job templates, respectively, when apply() is called on the parent job template.

TYPE: ChildJobTemplateLike DEFAULT: ()

consume_kwargs_from_results

Whether keyword arguments matched by a child job should be removed from the accumulated DAG results before later child jobs are matched.

TYPE: bool DEFAULT: True

name

Name of the job template. If not provided, the name of the wrapped callable is used.

TYPE: str | None DEFAULT: None

iterate_over_data_files

Whether dataset inputs should be processed item-wise.

TYPE: bool DEFAULT: False

output_dataset_param

Optional name of an explicit output-dataset parameter.

TYPE: str | None DEFAULT: None

output_dataset_cls

Optional dataset class to use for iterated outputs.

TYPE: type[IsDataset] | None DEFAULT: None

auto_async

Whether coroutine jobs at the outermost level (not in a flow context) should be automatically run in accordance with context (use existing event loop, if available, otherwise create temporary event loop and run coroutine until completion).

TYPE: bool DEFAULT: True

result_key

Optional key used to wrap the returned result in a dictionary. Especially useful in DAG flows to avoid name collisions.

TYPE: str | None DEFAULT: None

fixed_params

Fixed keyword-argument values for the job. May not target args or *kwargs-style params.

TYPE: Mapping[str, object] | Iterable[tuple[str, object]] | None DEFAULT: None

param_key_map

Mapping from callable parameter names to external keyword names. May not target args or *kwargs-style params.

TYPE: Mapping[str, str] | Iterable[tuple[str, str]] | None DEFAULT: None

persist_outputs

Per-job output-persistence preference.

TYPE: PersistOutputsOptions.Literals DEFAULT: PersistOutputsOptions.FOLLOW_CONFIG

restore_outputs

Per-job output-restore preference.

TYPE: RestoreOutputsOptions.Literals DEFAULT: RestoreOutputsOptions.FOLLOW_CONFIG

**kwargs

Additional constructor keyword overrides.

TYPE: object DEFAULT: {}

Returns: DagFlowTemplate: New DagFlowTemplate instance wrapping job_func.

Source code in src/omnipy/compute/flow.py
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
1353
1354
1355
1356
1357
1358
1359
1360
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
1373
1374
1375
1376
1377
1378
1379
1380
1381
1382
1383
1384
1385
1386
1387
1388
1389
1390
1391
1392
1393
1394
1395
1396
1397
1398
1399
1400
1401
1402
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
1464
1465
1466
1467
1468
1469
1470
1471
1472
1473
1474
1475
1476
1477
1478
1479
1480
1481
1482
1483
1484
1485
1486
1487
1488
1489
1490
1491
1492
1493
1494
1495
1496
1497
1498
1499
1500
1501
1502
1503
1504
1505
1506
1507
1508
1509
1510
1511
1512
1513
1514
1515
1516
1517
1518
1519
1520
1521
1522
1523
1524
1525
1526
1527
1528
1529
1530
1531
1532
1533
1534
1535
1536
1537
1538
1539
1540
1541
1542
1543
1544
1545
1546
1547
1548
1549
1550
1551
1552
1553
1554
1555
1556
1557
1558
1559
1560
1561
1562
1563
1564
1565
1566
1567
1568
1569
1570
1571
1572
1573
1574
1575
1576
1577
1578
1579
1580
1581
1582
1583
1584
1585
1586
1587
1588
1589
1590
1591
1592
1593
1594
1595
1596
1597
1598
1599
1600
1601
1602
1603
1604
1605
1606
1607
1608
1609
1610
1611
1612
1613
1614
1615
1616
1617
1618
1619
1620
1621
1622
1623
1624
1625
1626
1627
1628
1629
1630
1631
1632
1633
1634
1635
1636
1637
1638
1639
1640
1641
1642
1643
1644
1645
1646
1647
1648
1649
1650
1651
1652
1653
1654
1655
1656
1657
1658
1659
1660
1661
1662
1663
1664
1665
1666
1667
1668
1669
1670
1671
1672
1673
1674
1675
1676
1677
1678
1679
1680
1681
1682
1683
1684
1685
1686
1687
1688
1689
1690
1691
1692
1693
1694
1695
1696
1697
1698
1699
1700
1701
1702
1703
def DagFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    consume_kwargs_from_results: bool = True,
    name: str | None = None,
    iterate_over_data_files: bool = False,
    output_dataset_param: str | None = None,
    output_dataset_cls: type[IsDataset] | None = None,
    auto_async: bool = True,
    result_key: str | None = None,
    fixed_params: Mapping[str, object] | Iterable[tuple[str, object]] | None = None,
    param_key_map: Mapping[str, str] | Iterable[tuple[str, str]] | None = None,
    persist_outputs: PersistOutputsOptions.Literals = PersistOutputsOptions.FOLLOW_CONFIG,
    restore_outputs: RestoreOutputsOptions.Literals = RestoreOutputsOptions.FOLLOW_CONFIG,
    **kwargs: object,
) -> Any:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # Decorator-style factory for defining directed acyclic graph flows.
    #
    # {{DAG_FLOW_TEMPLATE_DESCRIPTION}}
    #
    # Args:
    #     {{JOB_TEMPLATE_CHILD_JOB_TEMPLATES_ARGS}}
    #     {{DAG_FLOW_TEMPLATE_KWARG_DOCS}}
    #     {{JOB_TEMPLATE_SHARED_KWARG_DOCS}}
    # Returns:
    #     DagFlowTemplate: New DagFlowTemplate instance wrapping ``job_func``.
    """Decorator-style factory for defining directed acyclic graph flows.

    A DAG flow template wraps a Python callable together with child job
    templates whose dependencies form a directed acyclic graph. Use this
    when flow steps branch and join but must not form cycles.

    ### Decorator usage

    Apply the template factory as a decorator to a Python callable.
    The wrapped callable becomes a reusable job template whose public outer
    signature is visible to template users and to the applied jobs created
    from it.

    ### Outer callable and child jobs

    The wrapped callable defines the public outer signature of the flow,
    while the child-job list defines the executed body.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True)
        ... def uppercase(data_file: TextModel) -> TextModel:
        ...     return data_file.content.upper()

        >>> @om.TaskTemplate()
        ... def join_texts(
        ...     upper: TextDataset,
        ...     original: TextDataset,
        ... ) -> TextDataset:
        ...     merged = TextDataset()
        ...     for title in upper:
        ...         merged[title] = f'{upper[title].content}|{original[title].content}'
        ...     return merged

        >>> @om.DagFlowTemplate(
        ...     uppercase.refine(result_key='upper'),
        ...     join_texts.refine(param_key_map={'upper': 'upper', 'original': 'dataset'}),
        ... )
        ... def my_dag(
        ...     dataset: TextDataset,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'HI|hi', 'b': 'BYE|bye'})
        >>> my_dag.run(text_files) == expected
        True

    ### Child jobs and data classes

    Child-job templates may be
    [TaskTemplate][omnipy.compute.task.TaskTemplate] or flow-template
    instances.
    This lets flows nest other flows as well as terminal tasks.

    In linear and DAG flows, Model and Dataset subclasses are also allowed
    as child-job entries. During apply, they are coerced into helper task
    templates that construct the requested data object from positional
    input in linear flows or from keyword-matched input in DAG flows.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.LinearFlowTemplate(TextModel)
        ... def wrap_text(raw_text: str) -> TextModel:
        ...     return TextModel(raw_text)

        >>> wrap_text.run('hello').content
        'hello'
        >>> @om.DagFlowTemplate(TextDataset)
        ... def collect_texts(first: str, second: str) -> TextDataset:
        ...     return TextDataset({'first': first, 'second': second})

        >>> collect_texts.run(first='hello', second='bye') == TextDataset({
        ...     'first': 'hello',
        ...     'second': 'bye',
        ... })
        True

    ### Callable-type validation

    For linear and DAG flows, the outer callable is primarily declarative:
    its signature exposes the public flow interface, while child jobs define
    the executed body.

    The outer callable type is validated against the child-job composition.
    The terminal child determines whether the flow behaves like a function
    or generator, and any async child lifts the full flow to an async
    callable type.

    As a result, a sync-function outer callable fits sync child execution,
    a generator outer callable fits generator-producing terminal children,
    and async outer callables are required when child composition is async.
    Mismatches raise ``TypeError`` when the flow template is created.

    ### Generator shorthand with ``Void``

    When the validated outer callable must be a generator or
    async-generator only to expose the correct public signature, use
    [Void][omnipy.compute.helpers.Void] in the body:

    Examples:
        >>> import omnipy as om
        >>> from collections.abc import Iterator

        >>> @om.TaskTemplate()
        ... def emit_lines() -> Iterator[str]:
        ...     yield 'first'
        ...     yield 'second'

        >>> @om.LinearFlowTemplate(emit_lines)
        ... def line_stream() -> Iterator[str]:
        ...     yield from om.Void()

    This shorthand exists only to satisfy the declared outer callable type;
    the child jobs still perform the actual flow work.

    ### Outer signature and modifiers

    The wrapped callable's parameter list and return annotation define the
    outer interface of the template.

    ``fixed_params`` permanently supplies selected callable parameters.

    ``param_key_map`` renames selected callable parameters to external
    keyword names that callers or parent flows use when supplying inputs.

    ``iterate_over_data_files``, ``output_dataset_param``, and
    ``output_dataset_cls`` adapt that outer interface for dataset-wise
    iteration.

    When ``iterate_over_data_files=True`` and the inner first parameter is
    annotated as ``Model[T]``, callers see an outer
    ``dataset: Dataset[Model[T]]`` parameter and the outer return type
    becomes a dataset of the per-item return type. The inner callable still
    receives one model object at a time.

    ``result_key`` wraps the returned value in a single-key dictionary,
    which is especially useful when a downstream DAG step should receive
    the result under a predictable name.

    Examples:
        >>> # With modifiers
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_other(number: int, other: int) -> int:
        ...     return number + other

        >>> plus_one = plus_other.refine(fixed_params={'other': 1})
        >>> plus_one.run(4)
        5
        >>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
        >>> plus_x.run(4, x=3)
        7
        >>> plus_one_dict = plus_one.refine(result_key='number')
        >>> plus_one_dict.run(4)
        {'number': 5}

    Examples:
        >>> # With dataset-wise iteration
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> add_suffix.run(text_files, suffix='!') == expected
        True

    ### DAG orchestration

    DAG flows route values by keyword name instead of chaining every child
    result positionally.

    The outer flow call first binds its arguments to the outer callable
    signature. Each child then receives the matching keyword subset from the
    accumulated named results.

    By default, a non-dictionary child result is stored under the child job
    name. ``result_key`` stores it under a custom key instead, which is the
    usual way to make one branch feed another. Dictionary results merge
    directly into the accumulated named result set.

    ``param_key_map`` lets a child read from externally visible DAG keys
    using different internal callable parameter names, and ``fixed_params``
    pins selected child inputs regardless of what earlier branches produce.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate()
        ... def uppercase(data_file: TextModel) -> TextModel:
        ...     return data_file.content.upper()

        >>> @om.TaskTemplate()
        ... def append_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> @om.TaskTemplate()
        ... def combine_texts(
        ...     left_dataset: TextDataset,
        ...     right_dataset: TextDataset,
        ... ) -> TextDataset:
        ...     merged = TextDataset()
        ...     for title in left_dataset:
        ...         merged[title] = (
        ...             f'{left_dataset[title].content}|'
        ...             f'{right_dataset[title].content}'
        ...         )
        ...     return merged

        >>> @om.DagFlowTemplate(
        ...     uppercase.refine(
        ...         iterate_over_data_files=True,
        ...         result_key='upper',
        ...     ),
        ...     append_suffix.refine(
        ...         iterate_over_data_files=True,
        ...         result_key='suffixed',
        ...         param_key_map={'suffix': 'ending'},
        ...     ),
        ...     combine_texts.refine(
        ...         param_key_map={
        ...             'left_dataset': 'upper',
        ...             'right_dataset': 'suffixed',
        ...         },
        ...     ),
        ... )
        ... def dag_flow(
        ...     dataset: TextDataset,
        ...     ending: str,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({
        ...     'a': 'HI|hi!',
        ...     'b': 'BYE|bye!',
        ... })
        >>> dag_flow.run(text_files, ending='!') == expected
        True

    ### Tasks and flows

    Tasks are terminal jobs: they wrap one callable and execute one compute
    step.

    Flows are orchestration jobs: they may contain child tasks and child
    flows, so larger pipelines can be assembled hierarchically from smaller
    reusable pieces.

    ### Lifecycle

    Apply a template with [`apply()`][omnipy.compute._job.JobTemplateMixin.apply]
    to create a runnable job with engine decorators and current config attached.
    Call the resulting applied job with runtime arguments.

    Use [`run()`][omnipy.compute._job.JobTemplateMixin.run] as a shorthand for
    ``apply()`` followed immediately by calling the applied job.

    Use [`refine()`][omnipy.compute._job.JobTemplateMixin.refine] to reuse a
    template while changing configuration such as ``name``, ``fixed_params``,
    or ``param_key_map``.

    Use [`revise()`][omnipy.compute._job.JobMixin.revise] on an applied job to
    reconstruct a template from that job's current configuration.

    Examples:
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_one(number: int) -> int:
        ...     return number + 1

        >>> plus_one.run(1)
        2
        >>> applied_job = plus_one.apply()
        >>> applied_job(2)
        3
        >>> refined_template = plus_one.refine(name='plus_one_renamed')
        >>> revised_template = applied_job.revise()

    Args:
        *child_job_templates: Ordered templates of child jobs to be
            run as part of the parent job. Model and Dataset subclasses
            are also allowed, in which case they are converted to
            create_X_from_args (if linear flow) or create_X_from_kwargs
            (if DAG flow) job templates, respectively, when apply() is
            called on the parent job template.
        consume_kwargs_from_results: Whether keyword arguments matched by a
            child job should be removed from the accumulated DAG results
            before later child jobs are matched.
        name: Name of the job template. If not provided, the name of the
            wrapped callable is used.
        iterate_over_data_files: Whether dataset inputs should be
            processed item-wise.
        output_dataset_param: Optional name of an explicit
            output-dataset parameter.
        output_dataset_cls: Optional dataset class to use for iterated
            outputs.
        auto_async: Whether coroutine jobs at the outermost level (not
            in a flow context) should be automatically run in accordance
            with context (use existing event loop, if available,
            otherwise create temporary event loop and run coroutine
            until completion).
        result_key: Optional key used to wrap the returned result in a
            dictionary. Especially useful in DAG flows to avoid name
            collisions.
        fixed_params: Fixed keyword-argument values for the job. May not
            target *args or **kwargs-style params.
        param_key_map: Mapping from callable parameter names to external
            keyword names. May not target *args or **kwargs-style
            params.
        persist_outputs: Per-job output-persistence preference.
        restore_outputs: Per-job output-restore preference.
        **kwargs: Additional constructor keyword overrides.
    Returns:
        DagFlowTemplate: New DagFlowTemplate instance wrapping ``job_func``."""
    ret = _DagFlowTemplateFactory(
        *child_job_templates,
        consume_kwargs_from_results=consume_kwargs_from_results,
        name=name,
        iterate_over_data_files=iterate_over_data_files,
        output_dataset_param=output_dataset_param,
        output_dataset_cls=output_dataset_cls,
        auto_async=auto_async,
        result_key=result_key,
        fixed_params=fixed_params,
        param_key_map=param_key_map,
        persist_outputs=persist_outputs,
        restore_outputs=restore_outputs,
        **kwargs,
    )
    return ret

FuncFlowTemplate

FuncFlowTemplate(
    *,
    iterate_over_data_files: Literal[True],
    output_dataset_cls: type[_RetDatasetClsT],
    **kwargs: Unpack[JobCommonKwargs],
) -> FuncFlowTemplateIterWithDatasetClsDecorator[_RetDatasetClsT]
FuncFlowTemplate(
    *,
    iterate_over_data_files: Literal[True],
    output_dataset_cls: None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> FuncFlowTemplateIterDecorator
FuncFlowTemplate(
    *,
    iterate_over_data_files: Literal[False] = False,
    output_dataset_cls: type[IsDataset] | None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> FuncFlowTemplatePlainDecorator
FuncFlowTemplate(
    *,
    iterate_over_data_files: bool,
    output_dataset_cls: type[_RetDatasetClsT] | None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> (
    FuncFlowTemplatePlainDecorator
    | FuncFlowTemplateIterDecorator
    | FuncFlowTemplateIterWithDatasetClsDecorator[_RetDatasetClsT]
)

Decorator-style factory for defining callable-backed coordinating flows.

A function flow template wraps a Python callable that orchestrates work as a flow. Use this when the control flow is easiest to express directly in Python instead of as an explicit task list or dependency graph.

Decorator usage

Apply the template factory as a decorator to a Python callable. The wrapped callable becomes a reusable job template whose public outer signature is visible to template users and to the applied jobs created from it.

Wrapped callable

The wrapped callable defines both the implementation and the public outer signature of the flow.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.FuncFlowTemplate()
... def append_suffix_to_all(
...     dataset: TextDataset,
...     suffix: str,
... ) -> TextDataset:
...     output_dataset = TextDataset()
...     for title, data_file in dataset.items():
...         output_dataset[title] = f'{data_file.content}{suffix}'
...     return output_dataset
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> append_suffix_to_all.run(text_files, '!') == expected
True
Outer signature and modifiers

The wrapped callable's parameter list and return annotation define the outer interface of the template.

fixed_params permanently supplies selected callable parameters.

param_key_map renames selected callable parameters to external keyword names that callers or parent flows use when supplying inputs.

iterate_over_data_files, output_dataset_param, and output_dataset_cls adapt that outer interface for dataset-wise iteration.

When iterate_over_data_files=True and the inner first parameter is annotated as Model[T], callers see an outer dataset: Dataset[Model[T]] parameter and the outer return type becomes a dataset of the per-item return type. The inner callable still receives one model object at a time.

result_key wraps the returned value in a single-key dictionary, which is especially useful when a downstream DAG step should receive the result under a predictable name.

Examples:

>>> # With modifiers
>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_other(number: int, other: int) -> int:
...     return number + other
>>> plus_one = plus_other.refine(fixed_params={'other': 1})
>>> plus_one.run(4)
5
>>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
>>> plus_x.run(4, x=3)
7
>>> plus_one_dict = plus_one.refine(result_key='number')
>>> plus_one_dict.run(4)
{'number': 5}

Examples:

>>> # With dataset-wise iteration
>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def add_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> add_suffix.run(text_files, suffix='!') == expected
True
Tasks and flows

Tasks are terminal jobs: they wrap one callable and execute one compute step.

Flows are orchestration jobs: they may contain child tasks and child flows, so larger pipelines can be assembled hierarchically from smaller reusable pieces.

Lifecycle

Apply a template with apply() to create a runnable job with engine decorators and current config attached. Call the resulting applied job with runtime arguments.

Use run() as a shorthand for apply() followed immediately by calling the applied job.

Use refine() to reuse a template while changing configuration such as name, fixed_params, or param_key_map.

Use revise() on an applied job to reconstruct a template from that job's current configuration.

Examples:

>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_one(number: int) -> int:
...     return number + 1
>>> plus_one.run(1)
2
>>> applied_job = plus_one.apply()
>>> applied_job(2)
3
>>> refined_template = plus_one.refine(name='plus_one_renamed')
>>> revised_template = applied_job.revise()
PARAMETER DESCRIPTION
name

Name of the job template. If not provided, the name of the wrapped callable is used.

TYPE: str | None DEFAULT: None

iterate_over_data_files

Whether dataset inputs should be processed item-wise.

TYPE: bool DEFAULT: False

output_dataset_param

Optional name of an explicit output-dataset parameter.

TYPE: str | None DEFAULT: None

output_dataset_cls

Optional dataset class to use for iterated outputs.

TYPE: type[IsDataset] | None DEFAULT: None

auto_async

Whether coroutine jobs at the outermost level (not in a flow context) should be automatically run in accordance with context (use existing event loop, if available, otherwise create temporary event loop and run coroutine until completion).

TYPE: bool DEFAULT: True

result_key

Optional key used to wrap the returned result in a dictionary. Especially useful in DAG flows to avoid name collisions.

TYPE: str | None DEFAULT: None

fixed_params

Fixed keyword-argument values for the job. May not target args or *kwargs-style params.

TYPE: Mapping[str, object] | Iterable[tuple[str, object]] | None DEFAULT: None

param_key_map

Mapping from callable parameter names to external keyword names. May not target args or *kwargs-style params.

TYPE: Mapping[str, str] | Iterable[tuple[str, str]] | None DEFAULT: None

persist_outputs

Per-job output-persistence preference.

TYPE: PersistOutputsOptions.Literals DEFAULT: PersistOutputsOptions.FOLLOW_CONFIG

restore_outputs

Per-job output-restore preference.

TYPE: RestoreOutputsOptions.Literals DEFAULT: RestoreOutputsOptions.FOLLOW_CONFIG

**kwargs

Additional constructor keyword overrides.

TYPE: object DEFAULT: {}

Returns: FuncFlowTemplate: New FuncFlowTemplate instance wrapping job_func.

Source code in src/omnipy/compute/flow.py
def FuncFlowTemplate(
    *,
    name: str | None = None,
    iterate_over_data_files: bool = False,
    output_dataset_param: str | None = None,
    output_dataset_cls: type[IsDataset] | None = None,
    auto_async: bool = True,
    result_key: str | None = None,
    fixed_params: Mapping[str, object] | Iterable[tuple[str, object]] | None = None,
    param_key_map: Mapping[str, str] | Iterable[tuple[str, str]] | None = None,
    persist_outputs: PersistOutputsOptions.Literals = PersistOutputsOptions.FOLLOW_CONFIG,
    restore_outputs: RestoreOutputsOptions.Literals = RestoreOutputsOptions.FOLLOW_CONFIG,
    **kwargs: object,
) -> Any:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # Decorator-style factory for defining callable-backed coordinating flows.
    #
    # {{FUNC_FLOW_TEMPLATE_DESCRIPTION}}
    #
    # Args:
    #     {{JOB_TEMPLATE_SHARED_KWARG_DOCS}}
    # Returns:
    #     FuncFlowTemplate: New FuncFlowTemplate instance wrapping ``job_func``.
    """Decorator-style factory for defining callable-backed coordinating flows.

    A function flow template wraps a Python callable that orchestrates
    work as a flow. Use this when the control flow is easiest to
    express directly in Python instead of as an explicit task list or
    dependency graph.

    ### Decorator usage

    Apply the template factory as a decorator to a Python callable.
    The wrapped callable becomes a reusable job template whose public outer
    signature is visible to template users and to the applied jobs created
    from it.

    ### Wrapped callable

    The wrapped callable defines both the implementation and the public
    outer signature of the flow.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.FuncFlowTemplate()
        ... def append_suffix_to_all(
        ...     dataset: TextDataset,
        ...     suffix: str,
        ... ) -> TextDataset:
        ...     output_dataset = TextDataset()
        ...     for title, data_file in dataset.items():
        ...         output_dataset[title] = f'{data_file.content}{suffix}'
        ...     return output_dataset

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> append_suffix_to_all.run(text_files, '!') == expected
        True

    ### Outer signature and modifiers

    The wrapped callable's parameter list and return annotation define the
    outer interface of the template.

    ``fixed_params`` permanently supplies selected callable parameters.

    ``param_key_map`` renames selected callable parameters to external
    keyword names that callers or parent flows use when supplying inputs.

    ``iterate_over_data_files``, ``output_dataset_param``, and
    ``output_dataset_cls`` adapt that outer interface for dataset-wise
    iteration.

    When ``iterate_over_data_files=True`` and the inner first parameter is
    annotated as ``Model[T]``, callers see an outer
    ``dataset: Dataset[Model[T]]`` parameter and the outer return type
    becomes a dataset of the per-item return type. The inner callable still
    receives one model object at a time.

    ``result_key`` wraps the returned value in a single-key dictionary,
    which is especially useful when a downstream DAG step should receive
    the result under a predictable name.

    Examples:
        >>> # With modifiers
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_other(number: int, other: int) -> int:
        ...     return number + other

        >>> plus_one = plus_other.refine(fixed_params={'other': 1})
        >>> plus_one.run(4)
        5
        >>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
        >>> plus_x.run(4, x=3)
        7
        >>> plus_one_dict = plus_one.refine(result_key='number')
        >>> plus_one_dict.run(4)
        {'number': 5}

    Examples:
        >>> # With dataset-wise iteration
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> add_suffix.run(text_files, suffix='!') == expected
        True

    ### Tasks and flows

    Tasks are terminal jobs: they wrap one callable and execute one compute
    step.

    Flows are orchestration jobs: they may contain child tasks and child
    flows, so larger pipelines can be assembled hierarchically from smaller
    reusable pieces.

    ### Lifecycle

    Apply a template with [`apply()`][omnipy.compute._job.JobTemplateMixin.apply]
    to create a runnable job with engine decorators and current config attached.
    Call the resulting applied job with runtime arguments.

    Use [`run()`][omnipy.compute._job.JobTemplateMixin.run] as a shorthand for
    ``apply()`` followed immediately by calling the applied job.

    Use [`refine()`][omnipy.compute._job.JobTemplateMixin.refine] to reuse a
    template while changing configuration such as ``name``, ``fixed_params``,
    or ``param_key_map``.

    Use [`revise()`][omnipy.compute._job.JobMixin.revise] on an applied job to
    reconstruct a template from that job's current configuration.

    Examples:
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_one(number: int) -> int:
        ...     return number + 1

        >>> plus_one.run(1)
        2
        >>> applied_job = plus_one.apply()
        >>> applied_job(2)
        3
        >>> refined_template = plus_one.refine(name='plus_one_renamed')
        >>> revised_template = applied_job.revise()

    Args:
        name: Name of the job template. If not provided, the name of the
            wrapped callable is used.
        iterate_over_data_files: Whether dataset inputs should be
            processed item-wise.
        output_dataset_param: Optional name of an explicit
            output-dataset parameter.
        output_dataset_cls: Optional dataset class to use for iterated
            outputs.
        auto_async: Whether coroutine jobs at the outermost level (not
            in a flow context) should be automatically run in accordance
            with context (use existing event loop, if available,
            otherwise create temporary event loop and run coroutine
            until completion).
        result_key: Optional key used to wrap the returned result in a
            dictionary. Especially useful in DAG flows to avoid name
            collisions.
        fixed_params: Fixed keyword-argument values for the job. May not
            target *args or **kwargs-style params.
        param_key_map: Mapping from callable parameter names to external
            keyword names. May not target *args or **kwargs-style
            params.
        persist_outputs: Per-job output-persistence preference.
        restore_outputs: Per-job output-restore preference.
        **kwargs: Additional constructor keyword overrides.
    Returns:
        FuncFlowTemplate: New FuncFlowTemplate instance wrapping ``job_func``."""
    ret = _FuncFlowTemplateFactory(
        name=name,
        iterate_over_data_files=iterate_over_data_files,
        output_dataset_param=output_dataset_param,
        output_dataset_cls=output_dataset_cls,
        auto_async=auto_async,
        result_key=result_key,
        fixed_params=fixed_params,
        param_key_map=param_key_map,
        persist_outputs=persist_outputs,
        restore_outputs=restore_outputs,
        **kwargs,
    )
    return ret

LinearFlowTemplate

LinearFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    iterate_over_data_files: Literal[True],
    output_dataset_cls: type[_RetDatasetClsT],
    **kwargs: Unpack[JobCommonKwargs],
) -> LinearFlowTemplateIterWithDatasetClsDecorator[_RetDatasetClsT]
LinearFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    iterate_over_data_files: Literal[True],
    output_dataset_cls: None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> LinearFlowTemplateIterDecorator
LinearFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    iterate_over_data_files: Literal[False] = False,
    output_dataset_cls: type[IsDataset] | None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> LinearFlowTemplatePlainDecorator
LinearFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    iterate_over_data_files: bool,
    output_dataset_cls: type[_RetDatasetClsT] | None = None,
    **kwargs: Unpack[JobCommonKwargs],
) -> (
    LinearFlowTemplatePlainDecorator
    | LinearFlowTemplateIterDecorator
    | LinearFlowTemplateIterWithDatasetClsDecorator[_RetDatasetClsT]
)

Decorator-style factory for defining sequential flows.

A linear flow template wraps a Python callable together with an ordered list of child job templates that run sequentially. Use this when work should proceed step by step in a fixed declaration order.

Decorator usage

Apply the template factory as a decorator to a Python callable. The wrapped callable becomes a reusable job template whose public outer signature is visible to template users and to the applied jobs created from it.

Outer callable and child jobs

The wrapped callable defines the public outer signature of the flow, while the child-job list defines the executed body.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def normalize_text(data_file: TextModel) -> TextModel:
...     return data_file.content.strip().lower()
>>> @om.LinearFlowTemplate(normalize_text)
... def clean_texts(
...     dataset: TextDataset,
... ) -> TextDataset:
...     return dataset
>>> text_files = TextDataset({'a': ' Hi ', 'b': 'BYE '})
>>> expected = TextDataset({'a': 'hi', 'b': 'bye'})
>>> clean_texts.run(text_files) == expected
True
Child jobs and data classes

Child-job templates may be TaskTemplate or flow-template instances. This lets flows nest other flows as well as terminal tasks.

In linear and DAG flows, Model and Dataset subclasses are also allowed as child-job entries. During apply, they are coerced into helper task templates that construct the requested data object from positional input in linear flows or from keyword-matched input in DAG flows.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.LinearFlowTemplate(TextModel)
... def wrap_text(raw_text: str) -> TextModel:
...     return TextModel(raw_text)
>>> wrap_text.run('hello').content
'hello'
>>> @om.DagFlowTemplate(TextDataset)
... def collect_texts(first: str, second: str) -> TextDataset:
...     return TextDataset({'first': first, 'second': second})
>>> collect_texts.run(first='hello', second='bye') == TextDataset({
...     'first': 'hello',
...     'second': 'bye',
... })
True
Callable-type validation

For linear and DAG flows, the outer callable is primarily declarative: its signature exposes the public flow interface, while child jobs define the executed body.

The outer callable type is validated against the child-job composition. The terminal child determines whether the flow behaves like a function or generator, and any async child lifts the full flow to an async callable type.

As a result, a sync-function outer callable fits sync child execution, a generator outer callable fits generator-producing terminal children, and async outer callables are required when child composition is async. Mismatches raise TypeError when the flow template is created.

Generator shorthand with Void

When the validated outer callable must be a generator or async-generator only to expose the correct public signature, use Void in the body:

Examples:

>>> import omnipy as om
>>> from collections.abc import Iterator
>>> @om.TaskTemplate()
... def emit_lines() -> Iterator[str]:
...     yield 'first'
...     yield 'second'
>>> @om.LinearFlowTemplate(emit_lines)
... def line_stream() -> Iterator[str]:
...     yield from om.Void()

This shorthand exists only to satisfy the declared outer callable type; the child jobs still perform the actual flow work.

Outer signature and modifiers

The wrapped callable's parameter list and return annotation define the outer interface of the template.

fixed_params permanently supplies selected callable parameters.

param_key_map renames selected callable parameters to external keyword names that callers or parent flows use when supplying inputs.

iterate_over_data_files, output_dataset_param, and output_dataset_cls adapt that outer interface for dataset-wise iteration.

When iterate_over_data_files=True and the inner first parameter is annotated as Model[T], callers see an outer dataset: Dataset[Model[T]] parameter and the outer return type becomes a dataset of the per-item return type. The inner callable still receives one model object at a time.

result_key wraps the returned value in a single-key dictionary, which is especially useful when a downstream DAG step should receive the result under a predictable name.

Examples:

>>> # With modifiers
>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_other(number: int, other: int) -> int:
...     return number + other
>>> plus_one = plus_other.refine(fixed_params={'other': 1})
>>> plus_one.run(4)
5
>>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
>>> plus_x.run(4, x=3)
7
>>> plus_one_dict = plus_one.refine(result_key='number')
>>> plus_one_dict.run(4)
{'number': 5}

Examples:

>>> # With dataset-wise iteration
>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def add_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> add_suffix.run(text_files, suffix='!') == expected
True
Linear orchestration

Linear flows run child jobs strictly in declaration order.

The first child receives the caller's positional and keyword inputs. Each later child receives the previous child result as its leading positional input, plus any matching keyword arguments from the outer flow call.

fixed_params always override caller-supplied values for the child where they are configured. param_key_map lets later children expose different external keyword names than their underlying callable parameter names.

Examples:

>>> import omnipy as om
>>> class TextModel(om.Model[str]): ...
>>> class TextDataset(om.Dataset[TextModel]): ...
>>> @om.TaskTemplate(iterate_over_data_files=True)
... def strip_text(data_file: TextModel) -> TextModel:
...     return data_file.content.strip()
>>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
... def add_suffix(
...     data_file: TextModel,
...     suffix: str,
... ) -> TextModel:
...     return f'{data_file.content}{suffix}'
>>> suffix_each = add_suffix.refine(param_key_map={'suffix': 'ending'})
>>> @om.LinearFlowTemplate(strip_text, suffix_each)
... def linear_flow(
...     dataset: TextDataset,
...     ending: str,
... ) -> TextDataset:
...     return dataset
>>> text_files = TextDataset({'a': ' hi', 'b': 'bye '})
>>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
>>> linear_flow.run(text_files, ending='!') == expected
True
Tasks and flows

Tasks are terminal jobs: they wrap one callable and execute one compute step.

Flows are orchestration jobs: they may contain child tasks and child flows, so larger pipelines can be assembled hierarchically from smaller reusable pieces.

Lifecycle

Apply a template with apply() to create a runnable job with engine decorators and current config attached. Call the resulting applied job with runtime arguments.

Use run() as a shorthand for apply() followed immediately by calling the applied job.

Use refine() to reuse a template while changing configuration such as name, fixed_params, or param_key_map.

Use revise() on an applied job to reconstruct a template from that job's current configuration.

Examples:

>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_one(number: int) -> int:
...     return number + 1
>>> plus_one.run(1)
2
>>> applied_job = plus_one.apply()
>>> applied_job(2)
3
>>> refined_template = plus_one.refine(name='plus_one_renamed')
>>> revised_template = applied_job.revise()
PARAMETER DESCRIPTION
*child_job_templates

Ordered templates of child jobs to be run as part of the parent job. Model and Dataset subclasses are also allowed, in which case they are converted to create_X_from_args (if linear flow) or create_X_from_kwargs (if DAG flow) job templates, respectively, when apply() is called on the parent job template.

TYPE: ChildJobTemplateLike DEFAULT: ()

name

Name of the job template. If not provided, the name of the wrapped callable is used.

TYPE: str | None DEFAULT: None

iterate_over_data_files

Whether dataset inputs should be processed item-wise.

TYPE: bool DEFAULT: False

output_dataset_param

Optional name of an explicit output-dataset parameter.

TYPE: str | None DEFAULT: None

output_dataset_cls

Optional dataset class to use for iterated outputs.

TYPE: type[IsDataset] | None DEFAULT: None

auto_async

Whether coroutine jobs at the outermost level (not in a flow context) should be automatically run in accordance with context (use existing event loop, if available, otherwise create temporary event loop and run coroutine until completion).

TYPE: bool DEFAULT: True

result_key

Optional key used to wrap the returned result in a dictionary. Especially useful in DAG flows to avoid name collisions.

TYPE: str | None DEFAULT: None

fixed_params

Fixed keyword-argument values for the job. May not target args or *kwargs-style params.

TYPE: Mapping[str, object] | Iterable[tuple[str, object]] | None DEFAULT: None

param_key_map

Mapping from callable parameter names to external keyword names. May not target args or *kwargs-style params.

TYPE: Mapping[str, str] | Iterable[tuple[str, str]] | None DEFAULT: None

persist_outputs

Per-job output-persistence preference.

TYPE: PersistOutputsOptions.Literals DEFAULT: PersistOutputsOptions.FOLLOW_CONFIG

restore_outputs

Per-job output-restore preference.

TYPE: RestoreOutputsOptions.Literals DEFAULT: RestoreOutputsOptions.FOLLOW_CONFIG

**kwargs

Additional constructor keyword overrides.

TYPE: object DEFAULT: {}

Returns: LinearFlowTemplate: New LinearFlowTemplate instance wrapping job_func.

Source code in src/omnipy/compute/flow.py
def LinearFlowTemplate(
    *child_job_templates: ChildJobTemplateLike,
    name: str | None = None,
    iterate_over_data_files: bool = False,
    output_dataset_param: str | None = None,
    output_dataset_cls: type[IsDataset] | None = None,
    auto_async: bool = True,
    result_key: str | None = None,
    fixed_params: Mapping[str, object] | Iterable[tuple[str, object]] | None = None,
    param_key_map: Mapping[str, str] | Iterable[tuple[str, str]] | None = None,
    persist_outputs: PersistOutputsOptions.Literals = PersistOutputsOptions.FOLLOW_CONFIG,
    restore_outputs: RestoreOutputsOptions.Literals = RestoreOutputsOptions.FOLLOW_CONFIG,
    **kwargs: object,
) -> Any:
    # %% Original docstring (managed by expand_docstr_macros.py) %%
    # Decorator-style factory for defining sequential flows.
    #
    # {{LINEAR_FLOW_TEMPLATE_DESCRIPTION}}
    #
    # Args:
    #     {{JOB_TEMPLATE_CHILD_JOB_TEMPLATES_ARGS}}
    #     {{JOB_TEMPLATE_SHARED_KWARG_DOCS}}
    # Returns:
    #     LinearFlowTemplate: New LinearFlowTemplate instance wrapping ``job_func``.
    """Decorator-style factory for defining sequential flows.

    A linear flow template wraps a Python callable together with an ordered
    list of child job templates that run sequentially. Use this when work
    should proceed step by step in a fixed declaration order.

    ### Decorator usage

    Apply the template factory as a decorator to a Python callable.
    The wrapped callable becomes a reusable job template whose public outer
    signature is visible to template users and to the applied jobs created
    from it.

    ### Outer callable and child jobs

    The wrapped callable defines the public outer signature of the flow,
    while the child-job list defines the executed body.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def normalize_text(data_file: TextModel) -> TextModel:
        ...     return data_file.content.strip().lower()

        >>> @om.LinearFlowTemplate(normalize_text)
        ... def clean_texts(
        ...     dataset: TextDataset,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': ' Hi ', 'b': 'BYE '})
        >>> expected = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> clean_texts.run(text_files) == expected
        True

    ### Child jobs and data classes

    Child-job templates may be
    [TaskTemplate][omnipy.compute.task.TaskTemplate] or flow-template
    instances.
    This lets flows nest other flows as well as terminal tasks.

    In linear and DAG flows, Model and Dataset subclasses are also allowed
    as child-job entries. During apply, they are coerced into helper task
    templates that construct the requested data object from positional
    input in linear flows or from keyword-matched input in DAG flows.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.LinearFlowTemplate(TextModel)
        ... def wrap_text(raw_text: str) -> TextModel:
        ...     return TextModel(raw_text)

        >>> wrap_text.run('hello').content
        'hello'
        >>> @om.DagFlowTemplate(TextDataset)
        ... def collect_texts(first: str, second: str) -> TextDataset:
        ...     return TextDataset({'first': first, 'second': second})

        >>> collect_texts.run(first='hello', second='bye') == TextDataset({
        ...     'first': 'hello',
        ...     'second': 'bye',
        ... })
        True

    ### Callable-type validation

    For linear and DAG flows, the outer callable is primarily declarative:
    its signature exposes the public flow interface, while child jobs define
    the executed body.

    The outer callable type is validated against the child-job composition.
    The terminal child determines whether the flow behaves like a function
    or generator, and any async child lifts the full flow to an async
    callable type.

    As a result, a sync-function outer callable fits sync child execution,
    a generator outer callable fits generator-producing terminal children,
    and async outer callables are required when child composition is async.
    Mismatches raise ``TypeError`` when the flow template is created.

    ### Generator shorthand with ``Void``

    When the validated outer callable must be a generator or
    async-generator only to expose the correct public signature, use
    [Void][omnipy.compute.helpers.Void] in the body:

    Examples:
        >>> import omnipy as om
        >>> from collections.abc import Iterator

        >>> @om.TaskTemplate()
        ... def emit_lines() -> Iterator[str]:
        ...     yield 'first'
        ...     yield 'second'

        >>> @om.LinearFlowTemplate(emit_lines)
        ... def line_stream() -> Iterator[str]:
        ...     yield from om.Void()

    This shorthand exists only to satisfy the declared outer callable type;
    the child jobs still perform the actual flow work.

    ### Outer signature and modifiers

    The wrapped callable's parameter list and return annotation define the
    outer interface of the template.

    ``fixed_params`` permanently supplies selected callable parameters.

    ``param_key_map`` renames selected callable parameters to external
    keyword names that callers or parent flows use when supplying inputs.

    ``iterate_over_data_files``, ``output_dataset_param``, and
    ``output_dataset_cls`` adapt that outer interface for dataset-wise
    iteration.

    When ``iterate_over_data_files=True`` and the inner first parameter is
    annotated as ``Model[T]``, callers see an outer
    ``dataset: Dataset[Model[T]]`` parameter and the outer return type
    becomes a dataset of the per-item return type. The inner callable still
    receives one model object at a time.

    ``result_key`` wraps the returned value in a single-key dictionary,
    which is especially useful when a downstream DAG step should receive
    the result under a predictable name.

    Examples:
        >>> # With modifiers
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_other(number: int, other: int) -> int:
        ...     return number + other

        >>> plus_one = plus_other.refine(fixed_params={'other': 1})
        >>> plus_one.run(4)
        5
        >>> plus_x = plus_other.refine(param_key_map={'other': 'x'})
        >>> plus_x.run(4, x=3)
        7
        >>> plus_one_dict = plus_one.refine(result_key='number')
        >>> plus_one_dict.run(4)
        {'number': 5}

    Examples:
        >>> # With dataset-wise iteration
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> text_files = TextDataset({'a': 'hi', 'b': 'bye'})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> add_suffix.run(text_files, suffix='!') == expected
        True

    ### Linear orchestration

    Linear flows run child jobs strictly in declaration order.

    The first child receives the caller's positional and keyword inputs.
    Each later child receives the previous child result as its leading
    positional input, plus any matching keyword arguments from the outer
    flow call.

    ``fixed_params`` always override caller-supplied values for the child
    where they are configured. ``param_key_map`` lets later children expose
    different external keyword names than their underlying callable
    parameter names.

    Examples:
        >>> import omnipy as om
        >>> class TextModel(om.Model[str]): ...
        >>> class TextDataset(om.Dataset[TextModel]): ...

        >>> @om.TaskTemplate(iterate_over_data_files=True)
        ... def strip_text(data_file: TextModel) -> TextModel:
        ...     return data_file.content.strip()

        >>> @om.TaskTemplate(iterate_over_data_files=True, output_dataset_cls=TextDataset)
        ... def add_suffix(
        ...     data_file: TextModel,
        ...     suffix: str,
        ... ) -> TextModel:
        ...     return f'{data_file.content}{suffix}'

        >>> suffix_each = add_suffix.refine(param_key_map={'suffix': 'ending'})
        >>> @om.LinearFlowTemplate(strip_text, suffix_each)
        ... def linear_flow(
        ...     dataset: TextDataset,
        ...     ending: str,
        ... ) -> TextDataset:
        ...     return dataset

        >>> text_files = TextDataset({'a': ' hi', 'b': 'bye '})
        >>> expected = TextDataset({'a': 'hi!', 'b': 'bye!'})
        >>> linear_flow.run(text_files, ending='!') == expected
        True

    ### Tasks and flows

    Tasks are terminal jobs: they wrap one callable and execute one compute
    step.

    Flows are orchestration jobs: they may contain child tasks and child
    flows, so larger pipelines can be assembled hierarchically from smaller
    reusable pieces.

    ### Lifecycle

    Apply a template with [`apply()`][omnipy.compute._job.JobTemplateMixin.apply]
    to create a runnable job with engine decorators and current config attached.
    Call the resulting applied job with runtime arguments.

    Use [`run()`][omnipy.compute._job.JobTemplateMixin.run] as a shorthand for
    ``apply()`` followed immediately by calling the applied job.

    Use [`refine()`][omnipy.compute._job.JobTemplateMixin.refine] to reuse a
    template while changing configuration such as ``name``, ``fixed_params``,
    or ``param_key_map``.

    Use [`revise()`][omnipy.compute._job.JobMixin.revise] on an applied job to
    reconstruct a template from that job's current configuration.

    Examples:
        >>> import omnipy as om
        >>> @om.TaskTemplate()
        ... def plus_one(number: int) -> int:
        ...     return number + 1

        >>> plus_one.run(1)
        2
        >>> applied_job = plus_one.apply()
        >>> applied_job(2)
        3
        >>> refined_template = plus_one.refine(name='plus_one_renamed')
        >>> revised_template = applied_job.revise()

    Args:
        *child_job_templates: Ordered templates of child jobs to be
            run as part of the parent job. Model and Dataset subclasses
            are also allowed, in which case they are converted to
            create_X_from_args (if linear flow) or create_X_from_kwargs
            (if DAG flow) job templates, respectively, when apply() is
            called on the parent job template.
        name: Name of the job template. If not provided, the name of the
            wrapped callable is used.
        iterate_over_data_files: Whether dataset inputs should be
            processed item-wise.
        output_dataset_param: Optional name of an explicit
            output-dataset parameter.
        output_dataset_cls: Optional dataset class to use for iterated
            outputs.
        auto_async: Whether coroutine jobs at the outermost level (not
            in a flow context) should be automatically run in accordance
            with context (use existing event loop, if available,
            otherwise create temporary event loop and run coroutine
            until completion).
        result_key: Optional key used to wrap the returned result in a
            dictionary. Especially useful in DAG flows to avoid name
            collisions.
        fixed_params: Fixed keyword-argument values for the job. May not
            target *args or **kwargs-style params.
        param_key_map: Mapping from callable parameter names to external
            keyword names. May not target *args or **kwargs-style
            params.
        persist_outputs: Per-job output-persistence preference.
        restore_outputs: Per-job output-restore preference.
        **kwargs: Additional constructor keyword overrides.
    Returns:
        LinearFlowTemplate: New LinearFlowTemplate instance wrapping ``job_func``."""
    ret = _LinearFlowTemplateFactory(
        *child_job_templates,
        name=name,
        iterate_over_data_files=iterate_over_data_files,
        output_dataset_param=output_dataset_param,
        output_dataset_cls=output_dataset_cls,
        auto_async=auto_async,
        result_key=result_key,
        fixed_params=fixed_params,
        param_key_map=param_key_map,
        persist_outputs=persist_outputs,
        restore_outputs=restore_outputs,
        **kwargs,
    )
    return ret