Skip to content

Commit ba43ae9

Browse files
chriscchienyangchiu
authored andcommitted
test(robot): add Test V2 Volume Engine Live Switchover
longhorn/longhorn#7125 Signed-off-by: Chris Chien <chris.chien@suse.com>
1 parent 600dbb7 commit ba43ae9

10 files changed

Lines changed: 351 additions & 9 deletions

File tree

e2e/keywords/workload.resource

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ Documentation Workload Keywords
33
44
Library Collections
55
Library String
6+
Library DateTime
67
Library ../libs/keywords/common_keywords.py
78
Library ../libs/keywords/volume_keywords.py
89
Library ../libs/keywords/workload_keywords.py
@@ -464,3 +465,55 @@ Check volume of ${workload_kind} ${workload_id} kept in faulted
464465
wait_for_volume_faulted ${volume_name}
465466
Sleep ${RETRY_INTERVAL}
466467
END
468+
469+
Start fio randwrite with crc32c verify in ${workload_kind} ${workload_id}
470+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
471+
start_fio_randwrite_with_verify_in_workload ${workload_name}
472+
473+
Stop fio in ${workload_kind} ${workload_id}
474+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
475+
stop_fio_in_workload ${workload_name}
476+
477+
Verify fio crc32c data in ${workload_kind} ${workload_id} have no errors
478+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
479+
verify_fio_data_integrity_in_workload ${workload_name}
480+
481+
Update volume of ${workload_kind} ${workload_id} engineNodeID to node ${node_id}
482+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
483+
${volume_name} = get_workload_volume_name ${workload_name}
484+
${node_name} = get_node_by_index ${node_id}
485+
update_volume_spec ${volume_name} engineNodeID ${node_name}
486+
487+
Volume of ${workload_kind} ${workload_id} should be attached to node ${node_id}
488+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
489+
${volume_name} = get_workload_volume_name ${workload_name}
490+
${expected_node_name} = get_node_by_index ${node_id}
491+
${expected_node_name} = get_node_by_index ${node_id}
492+
wait_for_volume_attached_to_node ${volume_name} ${expected_node_name}
493+
494+
Volume of ${workload_kind} ${workload_id} engine CR and enginefrontend CR should be on same node
495+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
496+
${volume_name} = get_workload_volume_name ${workload_name}
497+
wait_for_volume_engine_and_enginefrontend_same_node ${volume_name}
498+
499+
Volume of ${workload_kind} ${workload_id} engine CR should be on node ${node_id}
500+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
501+
${volume_name} = get_workload_volume_name ${workload_name}
502+
${expected_node_name} = get_node_by_index ${node_id}
503+
wait_for_volume_engine_node ${volume_name} ${expected_node_name}
504+
505+
Volume of ${workload_kind} ${workload_id} enginefrontend CR should be on node ${node_id}
506+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
507+
${volume_name} = get_workload_volume_name ${workload_name}
508+
${expected_node_name} = get_node_by_index ${node_id}
509+
wait_for_volume_enginefrontend_node ${volume_name} ${expected_node_name}
510+
511+
Mark volume monitoring start time for ${workload_kind} ${workload_id}
512+
${monitoring_start_time}= Get Current Date result_format=datetime
513+
Log Volume monitoring started at ${monitoring_start_time}
514+
Set Test Variable ${volume_monitoring_start_time} ${monitoring_start_time}
515+
516+
Volume of ${workload_kind} ${workload_id} should never have detached during test
517+
${workload_name} = generate_name_with_suffix ${workload_kind} ${workload_id}
518+
${volume_name} = get_workload_volume_name ${workload_name}
519+
verify_volume_never_detached_during_test ${volume_name} ${volume_monitoring_start_time}

e2e/libs/engine/engine.py

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,5 +59,19 @@ def get_engine_state(self, volume_name, node_name):
5959
def get_engine_name(self, volume_name):
6060
return self.engine.get_engine_name(volume_name)
6161

62+
def get_node(self, volume_name):
63+
for i in range(self.retry_count):
64+
logging(f"Trying to get engine node for volume {volume_name} ... ({i})")
65+
try:
66+
engine = self.get_engine(volume_name)
67+
engine_node = engine["spec"]["nodeID"]
68+
if engine_node:
69+
logging(f"Volume {volume_name} engine is running on node: {engine_node}")
70+
return engine_node
71+
except Exception as e:
72+
logging(f"Getting engine node for volume {volume_name} error: {e}")
73+
time.sleep(self.retry_interval)
74+
assert False, f"Failed to get engine node for volume {volume_name}"
75+
6276
def validate_engine_setting(self, volume_name, setting_name, value):
6377
return self.engine.validate_engine_setting(volume_name, setting_name, value)

e2e/libs/enginefrontend/enginefrontend.py

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,3 +43,22 @@ def get_enginefrontend_endpoint(self, volume_name):
4343
logging(f"Getting enginefrontend endpoint for volume {volume_name} error: {e}")
4444
time.sleep(self.retry_interval)
4545
assert False, f"Failed to get enginefrontend endpoint for volume {volume_name}"
46+
47+
def get_node(self, volume_name):
48+
"""
49+
Get the node where the volume's enginefrontend CR is running.
50+
"""
51+
for i in range(self.retry_count):
52+
logging(f"Trying to get enginefrontend node for volume {volume_name} ... ({i})")
53+
try:
54+
enginefrontend = self.get_enginefrontend(volume_name)
55+
enginefrontend_node = enginefrontend["spec"]["nodeID"]
56+
if enginefrontend_node:
57+
logging(f"Volume {volume_name} enginefrontend is running on node: {enginefrontend_node}")
58+
return enginefrontend_node
59+
else:
60+
raise RuntimeError(f"Unexpected empty enginefrontend nodeID: {enginefrontend['spec']}")
61+
except Exception as e:
62+
logging(f"Getting enginefrontend node for volume {volume_name} error: {e}")
63+
time.sleep(self.retry_interval)
64+
assert False, f"Failed to get enginefrontend node for volume {volume_name}"

e2e/libs/event/event.py

Lines changed: 34 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,11 +6,43 @@
66
from utility.utility import subprocess_exec_cmd
77

88

9-
def get_events(namespace=constant.LONGHORN_NAMESPACE):
9+
def get_events(namespace=constant.LONGHORN_NAMESPACE, field_selector=None, start_time=None):
10+
"""
11+
Get Kubernetes events from the specified namespace.
12+
13+
Args:
14+
namespace: Kubernetes namespace to query
15+
field_selector: Filter events by field selector (e.g., "involvedObject.name=volume-1,involvedObject.kind=Volume")
16+
start_time: Only return events after this datetime (must be timezone-aware or will be treated as UTC)
17+
18+
Returns:
19+
List of event items, optionally filtered by field_selector and start_time
20+
"""
1021
cmd = f"kubectl get event -n {namespace} --sort-by='.lastTimestamp' -ojson"
22+
23+
if field_selector:
24+
cmd += f" --field-selector='{field_selector}'"
25+
1126
try:
1227
events = json.loads(subprocess_exec_cmd(cmd, verbose=False))
13-
return events.get('items', [])
28+
items = events.get('items', [])
29+
30+
if start_time:
31+
if start_time.tzinfo is None:
32+
start_time = start_time.replace(tzinfo=timezone.utc)
33+
34+
filtered_items = []
35+
for event in items:
36+
# Parse lastTimestamp from the event
37+
last_timestamp_str = event.get('lastTimestamp')
38+
if last_timestamp_str:
39+
last_timestamp = datetime.fromisoformat(last_timestamp_str.replace('Z', '+00:00'))
40+
if last_timestamp >= start_time:
41+
filtered_items.append(event)
42+
43+
return filtered_items
44+
45+
return items
1446
except Exception as e:
1547
logging(f"Failed to get events in namespace {namespace}: {e}")
1648
return []

e2e/libs/keywords/volume_keywords.py

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
import asyncio
22
import time
33

4+
from engine import Engine
5+
from enginefrontend import EngineFrontend
46
from node import Node
57
from node.utility import check_replica_locality
68

@@ -22,6 +24,8 @@
2224
class volume_keywords:
2325

2426
def __init__(self):
27+
self.engine = Engine()
28+
self.enginefrontend = EngineFrontend()
2529
self.node = Node()
2630
self.volume = Volume()
2731
self.replica = Replica()
@@ -90,6 +94,37 @@ def wait_for_volume_not_attached_to_node(self, volume_name, node_name):
9094
time.sleep(self.retry_interval)
9195
assert False, f"Failed to wait for volume {volume_name} not attached to node {node_name}"
9296

97+
def wait_for_volume_engine_node(self, volume_name, expected_node_name):
98+
for i in range(self.retry_count):
99+
logging(f"Waiting for volume {volume_name} engine to be on node {expected_node_name} ... ({i})")
100+
engine_node = self.engine.get_node(volume_name)
101+
if engine_node == expected_node_name:
102+
return
103+
time.sleep(self.retry_interval)
104+
assert False, f"Failed to wait for volume {volume_name} engine on node {expected_node_name}, actual: {engine_node}"
105+
106+
def wait_for_volume_enginefrontend_node(self, volume_name, expected_node_name):
107+
for i in range(self.retry_count):
108+
logging(f"Waiting for volume {volume_name} enginefrontend to be on node {expected_node_name} ... ({i})")
109+
enginefrontend_node = self.enginefrontend.get_node(volume_name)
110+
if enginefrontend_node == expected_node_name:
111+
return
112+
time.sleep(self.retry_interval)
113+
assert False, f"Failed to wait for volume {volume_name} enginefrontend on node {expected_node_name}, actual: {enginefrontend_node}"
114+
115+
def wait_for_volume_engine_and_enginefrontend_same_node(self, volume_name):
116+
for i in range(self.retry_count):
117+
logging(f"Waiting for volume {volume_name} engine and enginefrontend to be on same node ... ({i})")
118+
volume_node = self.get_volume_node(volume_name)
119+
engine_node = self.engine.get_node(volume_name)
120+
enginefrontend_node = self.enginefrontend.get_node(volume_name)
121+
122+
if volume_node == engine_node and volume_node == enginefrontend_node:
123+
logging(f"Volume {volume_name} engine and enginefrontend are on same node: {volume_node}")
124+
return
125+
time.sleep(self.retry_interval)
126+
assert False, f"Failed to wait for volume {volume_name} engine and enginefrontend on same node. Volume node: {volume_node}, Engine node: {engine_node}, EngineFrontend node: {enginefrontend_node}"
127+
93128
def get_volume_instance_manager(self, volume_name):
94129
volume = VolumeRest().get(volume_name)
95130
assert len(volume['controllers']) == 1, f"Expect only one controller for volume {volume_name}; Got controllers: {volume['controllers']}"
@@ -478,3 +513,15 @@ def expand_volume(self, volume_name, size):
478513

479514
def trim_volume(self, volume_name):
480515
self.volume.trim_filesystem(volume_name)
516+
517+
def get_volume_engine_node(self, volume_name):
518+
return self.engine.get_node(volume_name)
519+
520+
def get_volume_enginefrontend_node(self, volume_name):
521+
return self.enginefrontend.get_node(volume_name)
522+
523+
def get_volume_state(self, volume_name):
524+
return self.volume.get_state(volume_name)
525+
526+
def verify_volume_never_detached_during_test(self, volume_name, start_time=None):
527+
return self.volume.verify_never_detached_during_test(volume_name, start_time)

e2e/libs/keywords/workload_keywords.py

Lines changed: 66 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2,14 +2,12 @@
22
import asyncio
33
import time
44

5-
from kubernetes import client
6-
from kubernetes.stream import stream
7-
85
from node import Node
96

107
from persistentvolumeclaim import PersistentVolumeClaim
118

129
from utility.utility import get_retry_count_and_interval
10+
from utility.utility import pod_exec
1311

1412
from workload.pod import get_volume_name_by_pod
1513
from workload.pod import new_busybox_manifest
@@ -422,3 +420,68 @@ def check_workload_pods_not_recreated(self, workload_name, workload_kind, namesp
422420
def rollout_restart_workload(self, workload_name, workload_kind, namespace="default"):
423421
logging(f"Triggering rollout restart of {workload_kind} {workload_name} in namespace {namespace}")
424422
rollout_restart_workload(workload_name, workload_kind, namespace)
423+
424+
def start_fio_randwrite_with_verify_in_workload(self, workload_name, namespace="default"):
425+
"""
426+
Start fio with randwrite and crc32c verification in workload pod.
427+
This runs fio in the background for data integrity testing.
428+
"""
429+
logging(f"Starting fio randwrite with crc32c verification in workload {workload_name}")
430+
pod_name = get_workload_pod_names(workload_name, namespace)[0]
431+
432+
fio_cmd = (
433+
"nohup fio --name=mytest "
434+
"--filename=/data/testfile.dat "
435+
"--size=1400M "
436+
"--rw=randwrite "
437+
"--bs=4k "
438+
"--ioengine=libaio "
439+
"--iodepth=32 "
440+
"--direct=1 "
441+
"--verify=crc32c "
442+
"--time_based "
443+
"--runtime=999999 "
444+
"--verify_state_save=1 "
445+
"> /data/fio_write.log 2>&1 &"
446+
)
447+
448+
resp = pod_exec(pod_name, namespace, fio_cmd)
449+
logging(f"Started fio in workload {workload_name}: {resp}")
450+
451+
def stop_fio_in_workload(self, workload_name, namespace="default"):
452+
"""
453+
Stop fio process in workload pod.
454+
"""
455+
logging(f"Stopping fio in workload {workload_name}")
456+
pod_name = get_workload_pod_names(workload_name, namespace)[0]
457+
458+
resp = pod_exec(pod_name, namespace, "pkill -SIGTERM fio || true")
459+
logging(f"Stopped fio in workload {workload_name}: {resp}")
460+
461+
def verify_fio_data_integrity_in_workload(self, workload_name, namespace="default"):
462+
"""
463+
Verify fio data integrity by running fio read with crc32c verification.
464+
Should return err=0 if data is intact.
465+
"""
466+
logging(f"Verifying fio data integrity in workload {workload_name}")
467+
pod_name = get_workload_pod_names(workload_name, namespace)[0]
468+
469+
verify_cmd = (
470+
"fio --name=mytest "
471+
"--filename=/data/testfile.dat "
472+
"--rw=read "
473+
"--bs=4k "
474+
"--ioengine=libaio "
475+
"--iodepth=32 "
476+
"--direct=1 "
477+
"--verify=crc32c "
478+
"--verify_state_load=1"
479+
)
480+
481+
resp = pod_exec(pod_name, namespace, verify_cmd)
482+
logging(f"FIO verification output: {resp}")
483+
484+
if "err= 0" not in resp and "err=0" not in resp:
485+
raise AssertionError(f"FIO verification failed with errors: {resp}")
486+
487+
logging(f"FIO data integrity verification passed")

e2e/libs/volume/crd.py

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@
1919
from volume.constant import GIBIBYTE, MEBIBYTE
2020
from volume.rest import Rest
2121

22+
from event.event import get_events
23+
2224

2325
class CRD(Base):
2426

@@ -886,3 +888,49 @@ def check_volume_has_recurringjob_group(self, volume_name, job_group_name):
886888
logging(f"Failed to find recurring job group {job_group_name} for volume {volume_name}")
887889
time.sleep(self.retry_interval)
888890
assert False, f"Failed to find recurring job group {job_group_name} for volume {volume_name}"
891+
892+
def get_state(self, volume_name):
893+
"""
894+
Get the current state of a volume.
895+
"""
896+
logging(f"Getting state for volume {volume_name}")
897+
volume = self.get(volume_name)
898+
state = volume["status"]["state"]
899+
logging(f"Volume {volume_name} state: {state}")
900+
return state
901+
902+
def verify_never_detached_during_test(self, volume_name, start_time=None):
903+
"""
904+
Verify that the volume never detached during the test by checking
905+
Kubernetes events for detachment events since start_time.
906+
907+
Args:
908+
volume_name: Name of the volume to check
909+
start_time: datetime object marking when to start checking events.
910+
If None, checks all events (not recommended).
911+
"""
912+
logging(f"Verifying volume {volume_name} never detached during test")
913+
914+
field_selector = f"involvedObject.name={volume_name},involvedObject.kind=Volume"
915+
events = get_events(
916+
namespace=constant.LONGHORN_NAMESPACE,
917+
field_selector=field_selector,
918+
start_time=start_time
919+
)
920+
921+
detachment_keywords = ["detached", "detaching", "disconnected"]
922+
923+
for event in events:
924+
event_message = event.get('message', '').lower()
925+
event_reason = event.get('reason', '').lower()
926+
event_timestamp = event.get('lastTimestamp', '')
927+
928+
for keyword in detachment_keywords:
929+
if keyword in event_message or keyword in event_reason:
930+
logging(f"Found detachment event: {event.get('reason')}: {event.get('message')} at {event_timestamp}")
931+
raise AssertionError(
932+
f"Volume {volume_name} had a detachment event during test: "
933+
f"{event.get('reason')}: {event.get('message')} at {event_timestamp}"
934+
)
935+
936+
logging(f"No detachment events found for volume {volume_name} since {start_time}")

e2e/libs/volume/volume.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -218,3 +218,9 @@ def check_volume_has_recurringjob(self, volume_name, job_name):
218218

219219
def check_volume_has_recurringjob_group(self, volume_name, job_group_name):
220220
self.volume.check_volume_has_recurringjob_group(volume_name, job_group_name)
221+
222+
def get_state(self, volume_name):
223+
return self.volume.get_state(volume_name)
224+
225+
def verify_never_detached_during_test(self, volume_name, start_time=None):
226+
return self.volume.verify_never_detached_during_test(volume_name, start_time)

e2e/libs/workload/deployment.py

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,15 @@ def create_deployment(name, claim_name, replicaset=1, enable_pvc_io_and_liveness
8282
# command is already set in the template, so we only need to set args
8383
manifest_dict['spec']['template']['spec']['containers'][0]['args'] = [args]
8484

85+
# If args contain apk commands, we need to run as root
86+
if 'apk' in args:
87+
# Override pod-level securityContext to run as root
88+
manifest_dict['spec']['template']['spec']['securityContext'] = {
89+
"runAsUser": 0,
90+
"runAsGroup": 0,
91+
"fsGroup": 0
92+
}
93+
8594
api = client.AppsV1Api()
8695

8796
deployment = api.create_namespaced_deployment(

0 commit comments

Comments
 (0)