forked from It4innovations/hyperqueue
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtest_job_mn.py
More file actions
131 lines (102 loc) · 4.75 KB
/
Copy pathtest_job_mn.py
File metadata and controls
131 lines (102 loc) · 4.75 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
import time
import pytest
from .conftest import HqEnv
from .utils import wait_for_job_state
from .utils.job import default_task_output
def test_submit_mn(hq_env: HqEnv):
hq_env.start_server()
hq_env.start_workers(2)
hq_env.command(["submit", "--nodes=3", "--", "bash", "-c", "sleep 1; echo ${HQ_NUM_NODES}; cat ${HQ_NODE_FILE}"])
time.sleep(0.5)
table = hq_env.command(["job", "info", "1"], as_table=True)
table.check_row_value("Resources", "nodes: 3")
table.check_row_value("State", "WAITING")
hq_env.start_workers(2)
wait_for_job_state(hq_env, 1, "RUNNING")
table = hq_env.command(["task", "list", "1"], as_table=True)
ws = table.get_column_value("Worker")[0].split("\n")
assert len(ws) == 3
assert set(ws).issubset(["worker1", "worker2", "worker3", "worker4"])
wait_for_job_state(hq_env, 1, "FINISHED", timeout_s=1.2)
with open(default_task_output(1)) as f:
assert f.readline() == "3\n"
hosts = f.read().rstrip().split("\n")
assert hosts == ws
# assert len(nodes) == 3
def test_reservation_in_mn(hq_env: HqEnv):
hq_env.start_server()
hq_env.command(["submit", "--nodes=3", "--", "bash", "-c", "sleep 3"])
workers = hq_env.start_workers(3, args=["--heartbeat=500ms", "--idle-timeout=1s"])
wait_for_job_state(hq_env, 1, "FINISHED")
hq_env.check_running_processes()
time.sleep(2)
for w in workers:
hq_env.check_process_exited(w, expected_code=None)
def test_failed_mn_task(hq_env: HqEnv):
hq_env.start_server()
hq_env.start_workers(3)
hq_env.command(["submit", "--nodes=3", "--", "bash", "-c", "exit 1"])
wait_for_job_state(hq_env, 1, "FAILED")
table = hq_env.command(["task", "list", "1"], as_table=True)
ws = table.get_column_value("Worker")[0].split("\n")
assert len(ws) == 3
assert set(ws).issubset(["worker1", "worker2", "worker3"])
def test_cancel_mn_task_running(hq_env: HqEnv):
hq_env.start_server()
hq_env.start_workers(3)
hq_env.command(["submit", "--nodes=3", "--", "bash", "-c", "sleep 10"])
wait_for_job_state(hq_env, 1, "RUNNING")
hq_env.command(["job", "cancel", "1"])
wait_for_job_state(hq_env, 1, "CANCELED")
hq_env.command(["submit", "--nodes=3", "--", "bash", "-c", "exit 0"])
wait_for_job_state(hq_env, 2, "FINISHED")
def test_cancel_mn_task_waiting(hq_env: HqEnv):
hq_env.start_server()
hq_env.start_worker()
hq_env.command(["submit", "--nodes=2", "--", "bash", "-c", "sleep 10"])
wait_for_job_state(hq_env, 1, "WAITING")
hq_env.command(["job", "cancel", "all"])
wait_for_job_state(hq_env, 1, "CANCELED")
@pytest.mark.parametrize("root", (True, False))
def test_worker_lost_mn_task(hq_env: HqEnv, root: bool):
hq_env.start_server()
hq_env.start_workers(3, cpus=1)
hq_env.command(["submit", "--nodes=3", "--", "bash", "-c", "sleep 2"])
wait_for_job_state(hq_env, 1, "RUNNING")
table = hq_env.command(["task", "list", "1"], as_table=True)
ws = table.get_column_value("Worker")[0].split("\n")
worker_ids = [int(w[len("worker") :]) for w in ws]
hq_env.kill_worker(worker_ids[0 if root else 1])
wait_for_job_state(hq_env, 1, "WAITING")
hq_env.start_workers(1, cpus=1)
wait_for_job_state(hq_env, 1, "FINISHED")
def test_submit_mn_different_groups(hq_env: HqEnv):
hq_env.start_server()
hq_env.start_workers(2, args=["--group=g1"])
hq_env.start_workers(2, args=["--group=g2"])
hq_env.command(["submit", "--nodes=3", "--", "/bin/hostname"])
time.sleep(0.5)
table = hq_env.command(["job", "info", "1"], as_table=True)
table.check_row_value("State", "WAITING")
hq_env.start_workers(1, args=["--group=g2"])
wait_for_job_state(hq_env, 1, "FINISHED")
def test_submit_mn_time_request(hq_env: HqEnv):
hq_env.start_server()
hq_env.start_workers(2)
hq_env.start_workers(1, args=["--time-limit=1s"], final_check=False)
hq_env.command(["submit", "--nodes=3", "--time-request=2s", "--", "/bin/hostname"])
time.sleep(0.7)
table = hq_env.command(["job", "info", "1"], as_table=True)
table.check_row_value("State", "WAITING")
hq_env.start_workers(1, args=["--time-limit=3s"])
wait_for_job_state(hq_env, 1, "FINISHED")
def test_submit_mn_complex_hostname(hq_env: HqEnv):
hq_env.start_server()
hq_env.start_worker(hostname="cn690.karolina.it4i.cz")
hq_env.start_worker(hostname="cn710.karolina.it4i.cz")
hq_env.command(["submit", "--nodes=2", "--", "bash", "-c", "sleep 1; cat ${HQ_NODE_FILE}; cat ${HQ_HOST_FILE}"])
wait_for_job_state(hq_env, 1, "FINISHED")
with open(default_task_output(1)) as f:
hosts = f.read().rstrip().split("\n")
assert sorted(hosts[:2]) == ["cn690", "cn710"]
assert sorted(hosts[2:4]) == ["cn690.karolina.it4i.cz", "cn710.karolina.it4i.cz"]