test_graph.py 4.84 KB
Newer Older
matthmey's avatar
matthmey committed
1
2
'''MIT License

matthmey's avatar
matthmey committed
3
4
Copyright (c) 2019, Swiss Federal Institute of Technology (ETH Zurich), Matthias Meyer

matthmey's avatar
matthmey committed
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23

Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:

The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.

THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.'''

24
import stuett
matthmey's avatar
matthmey committed
25
import datetime as dt
26
27
28
29
30
31
32
33
34
from tests.stuett.sample_data import *

import pytest

# helper function for non-stuett node
@stuett.dat
def bypass(x):
    return x

matthmey's avatar
matthmey committed
35

36
class MyNode(stuett.core.StuettNode):
matthmey's avatar
matthmey committed
37
38
39
40
    def __init__(self):
        super().__init__()

    def forward(self, data=None, request=None):
matthmey's avatar
matthmey committed
41
        print(data, request)
matthmey's avatar
matthmey committed
42
        return data + 4
43

matthmey's avatar
matthmey committed
44
    def configure(self, requests=None):
45
        requests = super().configure(requests)
matthmey's avatar
matthmey committed
46
47
        if "start_time" in requests:
            requests["start_time"] += 1
48
49
        return requests

matthmey's avatar
matthmey committed
50

matthmey's avatar
matthmey committed
51
52
53
class MyMerge(stuett.core.StuettNode):
    def forward(self, data, request):
        return data[0] + data[1]
54

matthmey's avatar
matthmey committed
55

56
class MySource(stuett.data.DataSource):
matthmey's avatar
matthmey committed
57
    def __init__(self, start_time=None):
matthmey's avatar
matthmey committed
58
59
        super().__init__(start_time=start_time)

matthmey's avatar
matthmey committed
60
    def forward(self, data=None, request=None):
matthmey's avatar
matthmey committed
61
62
        return request["start_time"]

63

matthmey's avatar
matthmey committed
64
65
66
67
def test_configuration():
    node = MyNode()

    # create a stuett graph
matthmey's avatar
matthmey committed
68
69
70
71
72
73
74
75
76
77
78
    x_in = node(data=5, delayed=True)  # data input
    x = bypass(x_in)
    x = node(x, delayed=True)

    print(x.compute())

    # x.visualize()
    import dask

    dsk, dsk_keys = dask.base._extract_graph_and_keys([x])
    print(dict(dsk))
matthmey's avatar
matthmey committed
79
80
81
82
83

    # create a configuration file
    config = {}

    # configure the graph
matthmey's avatar
matthmey committed
84
85
86
87
88
89
90
91
    x_configured = stuett.core.configuration(
        x, config
    )  # BUG: somehow config cuts the graph off

    import dask

    dsk, dsk_keys = dask.base._extract_graph_and_keys([x_configured])
    print(dsk)
matthmey's avatar
matthmey committed
92

matthmey's avatar
matthmey committed
93
94
95
    x = x_configured.compute()

    print(x)
matthmey's avatar
matthmey committed
96
97
    # TODO: finalize test

matthmey's avatar
matthmey committed
98
99
100
101

test_configuration()


matthmey's avatar
matthmey committed
102
103
104
def test_datasource():
    source = MySource()
    node = MyNode()
105

matthmey's avatar
matthmey committed
106
107
108
    # create a stuett graph
    x = source(delayed=True)
    x = bypass(x)
matthmey's avatar
matthmey committed
109
    x = node(x, delayed=True)
110

matthmey's avatar
matthmey committed
111
112
    # create a configuration file
    config = {"start_time": 0, "end_time": 1}
113

matthmey's avatar
matthmey committed
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
    # configure the graph
    configured = stuett.core.configuration(x, config)

    x_configured = configured.compute(
        scheduler="single-threaded", rerun_exceptions_locally=True
    )

    assert x_configured == 5


test_datasource()


def test_merging():
    source = MySource()
    node = MyNode()
    merge = MyMerge()
131

matthmey's avatar
matthmey committed
132
133
134
135
136
137
138
    # create a stuett graph
    import dask

    x = source(delayed=True)
    x = bypass(x)
    x = node(x, delayed=True)
    x_b = node(x, delayed=True)
matthmey's avatar
matthmey committed
139
    x = merge([x_b, x], delayed=True)
matthmey's avatar
matthmey committed
140
141
142

    # create a configuration file
    config = {"start_time": 0, "end_time": 1}
143

matthmey's avatar
matthmey committed
144
145
    # configure the graph
    configured = stuett.core.configuration(x, config)
146

matthmey's avatar
matthmey committed
147
    x_configured = configured.compute()
148

matthmey's avatar
matthmey committed
149
    assert x_configured == 14
150

matthmey's avatar
matthmey committed
151
    # TODO: Test default_merge
152

matthmey's avatar
matthmey committed
153
    # TODO:
154

matthmey's avatar
matthmey committed
155
    try:
matthmey's avatar
matthmey committed
156
        """
matthmey's avatar
matthmey committed
157
158
159
160
161
        This still fails. Configuration cannot handle "apply" nodes in the dask graph.
        The problem arises if we have a branch that needs to merge request and the 
        apply node is the one which receives the two requests. Then it requires a 
        default merge handler, which is not supplied (and does not need to) because the
        underlying function in this case supplies the merger.
matthmey's avatar
matthmey committed
162
        """
matthmey's avatar
matthmey committed
163
164
165
166
        x = source(delayed=True)(dask_key_name="source")
        x = bypass(x)(dask_key_name="bypass")
        x = node(x, delayed=True)(dask_key_name="node")
        x_b = node(x, delayed=True)(dask_key_name="branch")  # branch
matthmey's avatar
matthmey committed
167
        x = merge([x_b, x], delayed=True)(dask_key_name="merge")
168

matthmey's avatar
matthmey committed
169
170
        # create a configuration file
        config = {"start_time": 0, "end_time": 1}
171

matthmey's avatar
matthmey committed
172
        import dask
173

matthmey's avatar
matthmey committed
174
        dsk, dsk_keys = dask.base._extract_graph_and_keys([x])
175

matthmey's avatar
matthmey committed
176
        print(dict(dsk))
177

matthmey's avatar
matthmey committed
178
179
        # configure the graph
        configured = stuett.core.configuration(x, config)
180

matthmey's avatar
matthmey committed
181
        x_configured = configured.compute()
182

matthmey's avatar
matthmey committed
183
184
185
186
        assert x_configured == 14
    except Exception as e:
        print(e)
        pass
187

matthmey's avatar
matthmey committed
188
189

test_merging()