-
Notifications
You must be signed in to change notification settings - Fork 6
Expand file tree
/
Copy pathtest_backend.py
More file actions
151 lines (138 loc) · 5.63 KB
/
Copy pathtest_backend.py
File metadata and controls
151 lines (138 loc) · 5.63 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
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
from concurrent.futures import Future
import os
import shutil
import unittest
try:
from executorlib.task_scheduler.file.backend import backend_execute_task_in_file
from executorlib.task_scheduler.file.shared import _check_task_output, FutureItem
from executorlib.standalone.hdf import dump, get_runtime
from executorlib.standalone.serialize import serialize_funct
skip_h5io_test = False
except ImportError:
skip_h5io_test = True
def my_funct(a, b):
return a + b
def get_error(a):
raise ValueError(a)
@unittest.skipIf(
skip_h5io_test, "h5io is not installed, so the h5io tests are skipped."
)
class TestSharedFunctions(unittest.TestCase):
def test_execute_function_mixed(self):
cache_directory = os.path.abspath("executorlib_cache")
os.makedirs(cache_directory, exist_ok=True)
task_key, data_dict = serialize_funct(
fn=my_funct,
fn_args=[1],
fn_kwargs={"b": 2},
)
file_name = os.path.join(cache_directory, task_key + "_i.h5")
os.makedirs(cache_directory, exist_ok=True)
dump(file_name=file_name, data_dict=data_dict)
backend_execute_task_in_file(file_name=file_name)
future_obj = Future()
_check_task_output(
task_key=task_key, future_obj=future_obj, cache_directory=cache_directory
)
self.assertTrue(future_obj.done())
self.assertEqual(future_obj.result(), 3)
self.assertTrue(
get_runtime(file_name=os.path.join(cache_directory, task_key + "_o.h5"))
> 0.0
)
future_file_obj = FutureItem(
file_name=os.path.join(cache_directory, task_key + "_o.h5")
)
self.assertTrue(future_file_obj.done())
self.assertEqual(future_file_obj.result(), 3)
def test_execute_function_args(self):
cache_directory = os.path.abspath("executorlib_cache")
os.makedirs(cache_directory, exist_ok=True)
task_key, data_dict = serialize_funct(
fn=my_funct,
fn_args=[1, 2],
fn_kwargs=None,
)
file_name = os.path.join(cache_directory, task_key + "_i.h5")
os.makedirs(os.path.join(cache_directory, task_key), exist_ok=True)
dump(file_name=file_name, data_dict=data_dict)
backend_execute_task_in_file(file_name=file_name)
future_obj = Future()
_check_task_output(
task_key=task_key, future_obj=future_obj, cache_directory=cache_directory
)
self.assertTrue(future_obj.done())
self.assertEqual(future_obj.result(), 3)
self.assertTrue(
get_runtime(file_name=os.path.join(cache_directory, task_key + "_o.h5"))
> 0.0
)
future_file_obj = FutureItem(
file_name=os.path.join(cache_directory, task_key + "_o.h5")
)
self.assertTrue(future_file_obj.done())
self.assertEqual(future_file_obj.result(), 3)
def test_execute_function_kwargs(self):
cache_directory = os.path.abspath("executorlib_cache")
os.makedirs(cache_directory, exist_ok=True)
task_key, data_dict = serialize_funct(
fn=my_funct,
fn_args=None,
fn_kwargs={"a": 1, "b": 2},
)
file_name = os.path.join(cache_directory, task_key + "_i.h5")
os.makedirs(cache_directory, exist_ok=True)
dump(file_name=file_name, data_dict=data_dict)
backend_execute_task_in_file(file_name=file_name)
future_obj = Future()
_check_task_output(
task_key=task_key, future_obj=future_obj, cache_directory=cache_directory
)
self.assertTrue(future_obj.done())
self.assertEqual(future_obj.result(), 3)
self.assertTrue(
get_runtime(file_name=os.path.join(cache_directory, task_key + "_o.h5"))
> 0.0
)
future_file_obj = FutureItem(
file_name=os.path.join(cache_directory, task_key + "_o.h5")
)
self.assertTrue(future_file_obj.done())
self.assertEqual(future_file_obj.result(), 3)
def test_execute_function_error(self):
cache_directory = os.path.abspath("executorlib_cache")
os.makedirs(cache_directory, exist_ok=True)
task_key, data_dict = serialize_funct(
fn=get_error,
fn_args=[],
fn_kwargs={"a": 1},
)
file_name = os.path.join(cache_directory, task_key + "_i.h5")
os.makedirs(cache_directory, exist_ok=True)
data_dict["error_log_file"] = os.path.join(cache_directory, "error.out")
dump(file_name=file_name, data_dict=data_dict)
backend_execute_task_in_file(file_name=file_name)
future_obj = Future()
_check_task_output(
task_key=task_key, future_obj=future_obj, cache_directory=cache_directory
)
self.assertTrue(future_obj.done())
with self.assertRaises(ValueError):
future_obj.result()
with open(os.path.join(cache_directory, "error.out"), "r") as f:
content = f.readlines()
self.assertEqual(content[1], 'args: []\n')
self.assertEqual(content[2], "kwargs: {'a': 1}\n")
self.assertEqual(content[-1], 'ValueError: 1\n')
self.assertTrue(
get_runtime(file_name=os.path.join(cache_directory, task_key + "_o.h5"))
> 0.0
)
future_file_obj = FutureItem(
file_name=os.path.join(cache_directory, task_key + "_o.h5")
)
self.assertTrue(future_file_obj.done())
with self.assertRaises(ValueError):
future_file_obj.result()
def tearDown(self):
shutil.rmtree("executorlib_cache", ignore_errors=True)