-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathpipeline.py
More file actions
94 lines (76 loc) · 3.22 KB
/
Copy pathpipeline.py
File metadata and controls
94 lines (76 loc) · 3.22 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
# src/modacor/runner/pipeline.py
# # -*- coding: utf-8 -*-
from __future__ import annotations
from graphlib import TopologicalSorter
from pathlib import Path
import yaml
from attrs import define, field
from attrs import validators as v
from ..dataclasses.process_step import ProcessStep
__all__ = ["Pipeline"]
@define
class Pipeline(TopologicalSorter):
"""
Pipeline nodes are assumed to be of type ProcessStep
"""
name: str = field(default="Unnamed Pipeline")
graph: dict[ProcessStep] = field(factory=dict)
def __attrs_post_init__(self):
super().__init__(graph=self.graph)
@classmethod
def from_yaml(cls, path_to_yaml: Path):
"""
Instantiate a Pipeline from a yaml configuration file.
"""
yaml_obj = yaml.safe_load(path_to_yaml)
process_step_instances = {}
id_graph = {}
for module_name, module_data in yaml_obj["steps"].items():
# we need to instantiate ProcessSteps here, but
# need to implement a ProcessStep registry
step_id = module_data.get("step_id")
process_step_instances[step_id] = ProcessStep(io_sources=None)
id_graph[step_id] = module_data.get("requires_steps")
# translate step_id graph into ProcessStep graph
graph = {}
for k, v in id_graph.items():
graph[process_step_instances[k]] = {process_step_instances[i] for i in v}
return cls(name=yaml_obj["name"], graph=graph)
@classmethod
def from_dict(cls, graph_dict: dict, name=""):
return cls(name=name, graph=graph_dict)
def add_incoming_branch(self, branch: Self, branching_node):
"""
Add a pipeline as a branch whose outcome shall be combined the existing pipeline.
This assumes that the branch to be added has a single exit point.
"""
pipeline_to_add = Pipeline(graph=branch.graph)
pipeline_to_add_ordered = [*pipeline_to_add.static_order()]
# add the last node of the incoming as a predecessor to the connection point
self.graph[branching_node].update({pipeline_to_add_ordered[-1]})
# add the rest of the graph
self.graph = self.graph | branch.graph
# reinitialize the TopologicalSorter
super().__init__(graph=self.graph)
def add_outgoing_branch(self, branch: Self, branching_node):
"""
Add a pipeline as a branch whose input is based on the existing pipeline.
This assumes that the branch to be added has a single entry point.
"""
pipeline_to_add = Pipeline(graph=branch.graph)
pipeline_to_add_ordered = [*pipeline_to_add.static_order()]
# add the connection node as a predecessor to the first node of the outgoing branch
branch.graph[pipeline_to_add_ordered[0]].update({branching_node})
# add the rest of the graph
self.graph = self.graph | branch.graph
# reinitialize the TopologicalSorter
super().__init__(graph=self.graph)
def run(self, **kwargs):
"""
run pipeline with simple scheduling.
"""
self.prepare()
while self.is_active():
for node in self.get_ready():
node.execute(**kwargs)
self.done(node)