-
Notifications
You must be signed in to change notification settings - Fork 32
Expand file tree
/
Copy pathtest_liveliness.py
More file actions
326 lines (285 loc) · 11.6 KB
/
Copy pathtest_liveliness.py
File metadata and controls
326 lines (285 loc) · 11.6 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
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
"""Liveliness stress tests — parametrized over (sim × num_robots × iteration).
Verifies Docker stack bring-up: containers Running, sim publishes ``/clock``,
tmux panes host expected processes, sentinel ROS 2 nodes exist, compute
snapshots, and a short stability window (infra only — no camera/LiDAR Hz here).
Sensor topic rates, bridge stereo Hz, LiDAR echo/sanity, and sim RTF live in
``system/test_sensors.py`` (``@pytest.mark.sensors``), ordered after this module.
"""
import time
import pytest
from conftest import (
SimulatorHealthError,
collect_failure_diagnostics,
container_running,
current_test_id,
docker_exec,
find_all_containers,
get_metrics,
get_robot_containers,
logger,
ros2_exec,
sample_compute_usage,
wait_for_first_message,
)
SENTINEL_NODE_TEMPLATES = [
"/robot_{N}/interface/mavros/mavros",
"/robot_{N}/robot_state_publisher",
"/robot_{N}/trajectory_controller/trajectory_control_node",
]
def _parse_panes(raw):
"""Return (crashed, active_count). Input lines: 'session:window|pane_pid|title|kids'.
Crashed: pane with no direct children whose title isn't 'shell' (shell-tagged
panes are intentionally idle bash). Active: pane with at least one direct
child.
"""
crashed = []
active = 0
for line in raw.splitlines():
line = line.strip()
if not line:
continue
parts = line.split("|")
if len(parts) != 4:
continue
_, _, title, kids = parts
kid_count = int(kids.strip() or 0)
if kid_count == 0 and title != "shell":
crashed.append(line)
elif kid_count > 0:
active += 1
return crashed, active
def _check_tmux_panes(env):
"""Return (ok, msg). Pane with no direct child processes (and not a
'shell'-tagged idle bash) = crashed."""
cmd = (
"tmux list-panes -a -F '#{session_name}:#{window_name}|#{pane_pid}|#{pane_title}' "
"| while IFS='|' read -r w pid t; do "
"kids=$(pgrep -P \"$pid\" | wc -l); "
"printf '%s|%s|%s|%s\\n' \"$w\" \"$pid\" \"$t\" \"$kids\"; "
"done"
)
counts = {}
sim_container = env["sim_container"]
logger.info("Listing tmux panes in %s", sim_container)
result = docker_exec(sim_container, cmd, timeout=10)
if result.returncode != 0:
return False, f"tmux list-panes failed in {sim_container}"
sim_crashed, sim_active = _parse_panes(result.stdout)
counts[sim_container] = sim_active
if sim_crashed:
logger.warning("Sim panes crashed in %s: %s", sim_container, sim_crashed)
return False, f"sim panes crashed: {sim_crashed}"
for rc in get_robot_containers(env["robot_pattern"]):
logger.info("Listing tmux panes in %s", rc)
r = docker_exec(rc, cmd, timeout=10)
if r.returncode != 0:
return False, f"tmux list-panes failed in {rc}"
rcrashed, ractive = _parse_panes(r.stdout)
counts[rc] = ractive
if rcrashed:
logger.warning("Robot %s panes crashed: %s", rc, rcrashed)
return False, f"robot {rc} panes crashed: {rcrashed}"
summary = ", ".join(f"{c}={n}" for c, n in counts.items())
logger.info("All tmux panes active (%s)", summary)
return True, f"all tmux panes active ({summary})"
def _check_sim_startup_process(env):
"""Fast simulator-specific process/prerequisite health probe."""
if not container_running(env["sim_container"]):
return False, f"{env['sim_container']} stopped"
ok, message = _check_tmux_panes(env)
if not ok:
return ok, message
if env["sim"] != "msairsim":
return True, message
result = docker_exec(
env["sim_container"],
"binary=${MS_AIRSIM_BINARY_PATH:-"
"/ms-airsim-env/Blocks/LinuxNoEditor/Blocks.sh}; "
"test -x \"$binary\" && "
"nvidia-smi -L >/dev/null && "
"pgrep -fa 'Blocks|AirSim|UE4' >/dev/null",
timeout=10,
)
if result.returncode != 0:
return False, (
"Microsoft AirSim infrastructure prerequisite failed: scene binary "
"or GPU is unavailable, or the UE4 process exited"
)
return True, "Microsoft AirSim scene and UE4 process are healthy"
def _check_sentinel_nodes(env):
"""Return (ok, msg). Expected sentinels per robot domain."""
cfg = env["cfg"]
robot_containers = get_robot_containers(env["robot_pattern"])
if len(robot_containers) < env["num_robots"]:
return False, f"only {len(robot_containers)}/{env['num_robots']} robot containers visible"
all_missing = {}
for n in range(1, env["num_robots"] + 1):
logger.info(
"Checking sentinel nodes for robot_%d on domain %d in %s",
n,
n,
robot_containers[n - 1],
)
result = ros2_exec(
robot_containers[n - 1],
"ros2 node list 2>/dev/null",
domain_id=n,
setup_bash=cfg["robot_setup_bash"],
timeout=20,
)
if result.returncode != 0:
return False, f"ros2 node list failed for robot_{n}"
nodes = set(result.stdout.splitlines())
expected = {t.format(N=n) for t in SENTINEL_NODE_TEMPLATES}
missing = expected - nodes
if missing:
logger.warning("robot_%d missing nodes: %s", n, sorted(missing))
all_missing[f"robot_{n}"] = sorted(missing)
if all_missing:
return False, f"missing sentinel nodes: {all_missing}"
total = env["num_robots"] * len(SENTINEL_NODE_TEMPLATES)
logger.info("All %d sentinel nodes present", total)
return True, f"all {total} sentinel nodes present"
def _check_compute_usage(env):
"""Snapshot compute resources. Returns (ok, msg, samples_dict)."""
logger.info("Sampling compute usage")
try:
samples = sample_compute_usage(env["sim_container"])
except Exception as e:
logger.warning("Compute sampling raised: %s", e)
return False, f"compute sampling failed: {e}", {}
if not samples:
return False, "no compute samples returned", {}
logger.info("Sampled %d compute metrics", len(samples))
return True, f"{len(samples)} compute metrics sampled", samples
def _poll_until(predicate, timeout, interval, fail_msg):
"""Sleep-poll `predicate` up to `timeout` seconds."""
deadline = time.time() + timeout
while time.time() < deadline:
if predicate():
return
time.sleep(interval)
pytest.fail(fail_msg() if callable(fail_msg) else fail_msg)
@pytest.mark.liveliness
@pytest.mark.infrastructure
@pytest.mark.timeout(1800)
class TestLiveliness:
@pytest.mark.dependency(name="containers")
def test_robot_containers_running(self, airstack_env):
"""Wait up to 120s for N robot containers to be Running."""
num_robots = airstack_env["num_robots"]
pattern = airstack_env["robot_pattern"]
def ready():
containers = get_robot_containers(pattern)
return len(containers) >= num_robots and all(
container_running(c) for c in containers
)
_poll_until(
ready,
timeout=120,
interval=3,
fail_msg=lambda: f"only {len(get_robot_containers(pattern))}/"
f"{num_robots} robot containers Running after 120s",
)
@pytest.mark.dependency(name="sim_container", depends=["containers"])
def test_sim_container_running(self, airstack_env):
sc = airstack_env["sim_container"]
_poll_until(
lambda: container_running(sc),
timeout=120,
interval=3,
fail_msg=f"{sc} not Running after 120s",
)
@pytest.mark.dependency(depends=["containers"])
def test_gcs_container_running(self, airstack_env):
def ready():
names = find_all_containers("gcs")
return bool(names) and all(container_running(n) for n in names)
_poll_until(
ready,
timeout=120,
interval=3,
fail_msg="gcs container not Running after 120s",
)
@pytest.mark.dependency(name="sim_ready", depends=["sim_container"])
def test_sim_ready_time(self, airstack_env):
"""Wait for /clock while failing fast if the simulator process dies."""
cfg = airstack_env["cfg"]
m = get_metrics()
tid = current_test_id()
start = airstack_env["up_started_at"]
try:
ready = wait_for_first_message(
airstack_env["sim_container"],
"/clock",
domain_id=1,
setup_bash=cfg["sim_setup_bash"],
timeout=600,
health_check=lambda: _check_sim_startup_process(airstack_env),
health_grace=20,
)
except SimulatorHealthError as exc:
path = collect_failure_diagnostics(
airstack_env, str(exc), current_test_id()
)
pytest.fail(f"{exc}; diagnostics: {path}")
if ready is None:
m.record(tid, "sim_ready_duration_s", "timeout", unit="s")
path = collect_failure_diagnostics(
airstack_env,
"sim never published /clock within 600s",
current_test_id(),
)
pytest.fail(
f"sim never published /clock within 600s; diagnostics: {path}"
)
m.record(tid, "sim_ready_duration_s", round(time.time() - start, 2), unit="s")
@pytest.mark.dependency(name="tmux", depends=["containers"])
def test_tmux_panes_have_expected_processes(self, airstack_env):
ok, msg = _check_tmux_panes(airstack_env)
assert ok, msg
@pytest.mark.dependency(name="compute", depends=["sim_ready"])
def test_compute_usage(self, airstack_env):
"""Snapshot per-container CPU/mem/IO + host + GPU (diagnostic metrics)."""
ok, msg, _ = _check_compute_usage(airstack_env)
assert ok, msg
@pytest.mark.dependency(name="nodes", depends=["containers"])
def test_sentinel_nodes_present(self, airstack_env):
"""Wait up to 300s for the expected sentinel nodes per robot."""
last_msg = [""]
def ready():
ok, msg = _check_sentinel_nodes(airstack_env)
last_msg[0] = msg
return ok
_poll_until(
ready,
timeout=300,
interval=5,
fail_msg=lambda: f"sentinel nodes not ready after 300s: {last_msg[0]}",
)
@pytest.mark.dependency(depends=["sim_ready", "nodes", "tmux"])
def test_stable(self, airstack_env, request):
"""Poll infra only: tmux, sentinel nodes, compute (no sensor topic Hz)."""
duration = request.config.getoption("--stable-duration")
interval = request.config.getoption("--stable-interval")
m = get_metrics()
tid = current_test_id()
series = {}
elapsed = 0
try:
while elapsed < duration:
time.sleep(interval)
elapsed += interval
ok_t, msg_t = _check_tmux_panes(airstack_env)
ok_n, msg_n = _check_sentinel_nodes(airstack_env)
_, _, compute = _check_compute_usage(airstack_env)
for key, value in compute.items():
series.setdefault(key, []).append({"t": elapsed, "value": value})
if not (ok_t and ok_n):
pytest.fail(
f"instability at t={elapsed}s: tmux={msg_t} | nodes={msg_n}"
)
finally:
for key, samples in series.items():
if samples:
m.record_list(tid, f"{key}_samples", samples)