-
Notifications
You must be signed in to change notification settings - Fork 10
Expand file tree
/
Copy pathstorage.py
More file actions
491 lines (412 loc) · 19.2 KB
/
Copy pathstorage.py
File metadata and controls
491 lines (412 loc) · 19.2 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
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
from __future__ import annotations
import logging
from dataclasses import dataclass
from lib import config
from lib.commands import SSHCommandFailed
from lib.common import QCOW2_MAX, VHD_MAX, Defer, MiB, PackageManagerEnum, strtobool, wait_for
from lib.host import Host
from lib.sr import SR
from lib.vdi import VDI, ImageFormat
from lib.vm import VM
from typing import Literal
MAX_VDI_SIZE: dict[ImageFormat, int] = {'qcow2': QCOW2_MAX, 'vhd': VHD_MAX}
def try_to_create_sr_with_missing_device(sr_type, label, host) -> None:
try:
host.sr_create(sr_type, label, {}, verify=True)
except SSHCommandFailed as e:
assert e.stdout == (
'Error code: SR_BACKEND_FAILURE_90\nError parameters: , '
+ 'The request is missing the device parameter,'
), 'Bad error, current: {}'.format(e.stdout)
return
assert False, 'SR creation should not have succeeded!'
def cold_migration_then_come_back(vm: VM, prov_host: Host, dest_host: Host, dest_sr: SR) -> None:
""" Storage migration of a shutdown VM, then migrate it back. """
prov_sr = vm.get_sr()
vdi_name: str | None = None
integrity_check = not vm.is_windows
dev = ""
spans: list[StreamSpan] = []
if integrity_check:
# the vdi will be destroyed with the vm
vdi = prov_sr.create_vdi(virtual_size=config.volume_size)
vdi_name = vdi.name()
vbd = vm.connect_vdi(vdi)
vm.start()
vm.wait_for_vm_running_and_ssh_up()
install_randstream(vm)
dev = f'/dev/{vbd.param_get("device")}'
spans = partially_populate_device(vm, dev, config.volume_size)
validate_partially_populated_device(vm, dev, spans)
vm.shutdown(verify=True)
assert vm.is_halted()
# Move the VM to another host of the pool
vm.migrate(dest_host, dest_sr)
wait_for(lambda: vm.all_vdis_on_sr(dest_sr), "Wait for all VDIs on destination SR")
# Start VM to make sure it works
vm.start(on=dest_host.uuid)
vm.wait_for_os_booted()
if integrity_check:
vm.wait_for_vm_running_and_ssh_up()
validate_partially_populated_device(vm, dev, spans)
vm.shutdown(verify=True)
# Migrate it back to the provenance SR
vm.migrate(prov_host, prov_sr)
wait_for(lambda: vm.all_vdis_on_sr(prov_sr), "Wait for all VDIs back on provenance SR")
# Start VM to make sure it works
vm.start(on=prov_host.uuid)
vm.wait_for_os_booted()
if integrity_check:
vm.wait_for_vm_running_and_ssh_up()
validate_partially_populated_device(vm, dev, spans)
vm.shutdown(verify=True)
if vdi_name is not None:
vm.destroy_vdi_by_name(vdi_name)
def live_storage_migration_then_come_back(vm: VM, prov_host: Host, dest_host: Host, dest_sr: SR) -> None:
prov_sr = vm.get_sr()
vdi_name: str | None = None
integrity_check = not vm.is_windows
dev = ""
spans: list[StreamSpan] = []
vbd = None
if integrity_check:
vdi = prov_sr.create_vdi(virtual_size=config.volume_size)
vdi_name = vdi.name()
vbd = vm.connect_vdi(vdi)
# start VM
vm.start(on=prov_host.uuid)
vm.wait_for_os_booted()
if integrity_check:
vm.wait_for_vm_running_and_ssh_up()
install_randstream(vm)
assert vbd is not None
dev = f'/dev/{vbd.param_get("device")}'
spans = partially_populate_device(vm, dev, config.volume_size)
validate_partially_populated_device(vm, dev, spans)
# Move the VM to another host of the pool
vm.migrate(dest_host, dest_sr)
wait_for(lambda: vm.all_vdis_on_sr(dest_sr), "Wait for all VDIs on destination SR")
wait_for(lambda: vm.is_running_on_host(dest_host), "Wait for VM to be running on destination host")
if integrity_check:
validate_partially_populated_device(vm, dev, spans)
# Migrate it back to the provenance SR
vm.migrate(prov_host, prov_sr)
wait_for(lambda: vm.all_vdis_on_sr(prov_sr), "Wait for all VDIs back on provenance SR")
wait_for(lambda: vm.is_running_on_host(prov_host), "Wait for VM to be running on provenance host")
if integrity_check:
validate_partially_populated_device(vm, dev, spans)
vm.shutdown(verify=True)
if vdi_name is not None:
vm.destroy_vdi_by_name(vdi_name)
def vdi_is_open(vdi: VDI) -> bool:
sr = vdi.sr
get_sr_ref = f"""
import sys
import XenAPI
def get_xapi_session():
session = XenAPI.xapi_local()
try:
session.xenapi.login_with_password('root', '', '', 'xcp-ng-tests session')
except Exception as e:
raise Exception('Cannot get XAPI session: {{}}'.format(e))
return session
session = get_xapi_session()
try:
sr_ref = session.xenapi.SR.get_by_uuid(\"{sr.uuid}\")
finally:
session.xenapi.logout()
print(sr_ref)
"""
master = sr.pool.master
return strtobool(master.call_plugin('on-slave', 'is_open', {
'vdiUuid': vdi.uuid,
'srRef': master.execute_script(get_sr_ref, shebang='python')
}))
def install_randstream(vm: VM) -> None:
BASE_URL = 'https://github.com/xcp-ng/randstream/releases/download'
VERSION = '0.6.1'
CHECKSUM = {
'Linux': '2aef357cfdfed09d6492cfb60ab145f304abd7ccda9b77d47e15a9048c1a6eee',
'FreeBSD': '645c007393be939d75d95388b1d0e41e0e1217e8a88057284916b638a61c690d',
}
TARGET_TRIPLE = {
'Linux': 'x86_64-unknown-linux-musl',
'FreeBSD': 'x86_64-unknown-freebsd',
}
version = vm.ssh('randstream --version', check=False)
if f'randstream {VERSION}' == version:
logging.debug("randstream is already installed")
return
logging.debug("Installing randstream")
if vm.is_windows:
raise ValueError("Windows is not currently supported")
else:
os_name = vm.ssh('uname -s')
assert os_name in CHECKSUM, f"{os_name} is not currently supported"
tt = TARGET_TRIPLE[os_name]
cs = CHECKSUM[os_name]
fn = '/tmp/randstream.tgz'
vm.ssh(f"echo '{cs} -' > {fn}.sum && wget -nv {BASE_URL}/{VERSION}/randstream-{VERSION}-{tt}.tar.gz -O - | tee {fn} | sha256sum -c {fn}.sum && tar -xzf {fn} -C /usr/bin/ ./randstream") # noqa: E501
vm.ssh(f"rm -f {fn} {fn}.sum")
def randstream(vm: VM, args: str) -> str:
"""
Run randstream on the VM and return the checksum.
The args string should contain the command and arguments to pass to randstream,
e.g. "generate /dev/xvdb".
"""
output = vm.ssh(f'randstream -v {args}')
for line in output.splitlines():
if line.startswith('checksum: '):
return line.split(": ")[1].strip()
raise Exception(f"Could not find the checksum in the randstream output:\n{output}")
CoalesceOperation = Literal['snapshot', 'clone']
def coalesce_integrity(vm: VM, vdi: VDI, vdi_op: CoalesceOperation, defer: Defer) -> None:
vdi_size = vdi.get_virtual_size()
vbd = vm.connect_vdi(vdi)
defer(lambda: vm.disconnect_vdi(vdi))
dev = f'/dev/{vbd.param_get("device")}'
# generate at the start, in the middle and the end of the disk
spans = partially_populate_device(vm, dev, vdi_size, 4, skip_spans=[1])
# make sure we can read that exact data before the snapshot/clone
validate_partially_populated_device(vm, dev, spans)
new_vdi: VDI | None = None
match vdi_op:
case 'clone': new_vdi = vdi.clone()
case 'snapshot': new_vdi = vdi.snapshot()
defer(lambda: new_vdi.destroy() if new_vdi is not None else None)
assert vdi is not None
# add some data in a non-used place (span 1), and overwrite an already used one (span 2)
spans[1].generate(vm, dev, seed=1)
spans[2].generate(vm, dev, seed=2)
# make sure we can validate that data before the coalesce
spans[1].validate(vm, dev)
spans[2].validate(vm, dev)
# trigger the coalesce
vdi.wait_for_coalesce(new_vdi.destroy)
new_vdi = None
# verify the data is still as expected
validate_partially_populated_device(vm, dev, spans)
XVACompression = Literal['none', 'gzip', 'zstd']
def xva_export_import(source_vm: VM, compression: XVACompression, temp_large_dir: str, defer: Defer) -> None:
# clone the vm, so we can resize the disk without affecting the vm from the fixture
vm: VM | None = source_vm.clone()
defer(lambda: vm.destroy() if vm is not None else None)
assert vm is not None
host = vm.host
sr = vm.vdis[0].sr
# we can't shrink a volume
volume_size = max(vm.vdis[0].get_virtual_size(), config.volume_size)
vm.vdis[0].resize(volume_size)
# The resulting volume size is a multiple of the block size. Store the actual VDI size, so we can make comparisons
# later in the test
volume_size = vm.vdis[0].get_virtual_size()
# The tests using this function are using specific fixtures to create the VM on the expected SR
# In consequence, we can't use the storage_test_vm, so we have to start the VM explicitly and install randstream
vm.start()
vm.wait_for_vm_running_and_ssh_up()
install_randstream(vm)
if vm.detect_package_manager() == PackageManagerEnum.APK:
# growpart is not available in alpine 3.12
# vm.ssh('apk add cloud-utils-growpart e2fsprogs-extra')
vm.ssh('apk add gawk util-linux e2fsprogs-extra')
vm.ssh('wget https://raw.githubusercontent.com/canonical/cloud-utils/main/bin/growpart -O /usr/bin/growpart')
vm.ssh('chmod +x /usr/bin/growpart')
# TODO: maybe use `findmnt -no SOURCE /` from util-linux to get the blockdevice mounted on /
growpart_returncode = vm.ssh_with_result('growpart /dev/xvda 3').returncode
assert growpart_returncode in [0, 1] # growpart returns 1 if the size is already the expected one
vm.ssh('resize2fs /dev/xvda3')
stream_size = min(volume_size // 2, config.write_volume_cap)
else:
stream_size = 500 * MiB
checksum = randstream(vm, f'generate --size {stream_size} /root/data')
randstream(vm, f'validate --expected-checksum {checksum} /root/data')
vm.shutdown(verify=True)
xva_path = f'{temp_large_dir}/{vm.uuid}.xva'
defer(lambda: host.ssh(f'rm -f {xva_path}'))
vm.export(xva_path, compression)
# check that the zero blocks are not part of the result. Most of the data is from the random stream, so
# compression has little effect. We just take into account the system size
size_mb = int(vm.host.ssh(f'du -sm --apparent-size {xva_path}').split()[0])
min_size = stream_size / MiB
max_size = (stream_size + volume_size / 1000) * 1.1 / MiB + 200
assert min_size < size_mb < max_size, (
f"unexpected xva size {size_mb}MiB, was expected to be between {min_size}MiB and {max_size}MiB"
)
# destroy the source vm to free some space to re-import the image
vm.destroy()
vm = None
imported_vm = host.import_vm(xva_path, sr.uuid)
defer(lambda: imported_vm.destroy())
assert imported_vm.vdis[0].get_virtual_size() == volume_size
imported_vm.start()
imported_vm.wait_for_vm_running_and_ssh_up()
randstream(imported_vm, f'validate --expected-checksum {checksum} /root/data')
def vdi_export_import(vm: VM, sr: SR, image_format: ImageFormat, temp_large_dir: str, defer: Defer) -> None:
vdi_src: VDI | None = sr.create_vdi(image_format=image_format, virtual_size=config.volume_size)
defer(lambda: vdi_src.destroy() if vdi_src is not None else None)
assert vdi_src is not None
vbd = vm.connect_vdi(vdi_src)
defer(lambda: vm.disconnect_vdi(vdi_src) if vdi_src is not None and vdi_src.uuid in vm.vdis else None)
dev = f'/dev/{vbd.param_get("device")}'
spans = partially_populate_device(vm, dev, config.volume_size)
validate_partially_populated_device(vm, dev, spans)
vm.disconnect_vdi(vdi_src)
image_path = f'{temp_large_dir}/{vdi_src.uuid}.{image_format}'
defer(lambda: vm.host.ssh(f'rm -f {image_path}'))
vm.host.xe('vdi-export', {'uuid': vdi_src.uuid, 'filename': image_path, 'format': image_format})
vdi_src.destroy()
vdi_src = None
# check that the zero blocks are not part of the result
size_mb = int(vm.host.ssh(f'du -sm --apparent-size {image_path}').split()[0])
total_span_size_mib = sum(span.size for span in spans) // MiB
assert total_span_size_mib < size_mb < total_span_size_mib * 1.1, f"unexpected image size: {size_mb}"
vdi_dest = sr.create_vdi(image_format=image_format, virtual_size=config.volume_size)
defer(lambda: vdi_dest.destroy())
vm.host.xe('vdi-import', {'uuid': vdi_dest.uuid, 'filename': image_path, 'format': image_format})
vbd = vm.connect_vdi(vdi_dest)
defer(lambda: vm.disconnect_vdi(vdi_dest))
dev = f'/dev/{vbd.param_get("device")}'
validate_partially_populated_device(vm, dev, spans)
def full_vdi_write(vm: VM, vdi: VDI, defer: Defer):
vdi.get_virtual_size()
vbd = vm.connect_vdi(vdi)
defer(lambda: vm.disconnect_vdi(vdi))
dev = f'/dev/{vbd.param_get("device")}'
install_randstream(vm)
checksum = randstream(vm, f'generate {dev}')
randstream(vm, f'validate --expected-checksum {checksum} {dev}')
@dataclass
class StreamSpan:
position: int
size: int
checksum: str | None = None
def generate(self, vm: VM, dev: str, seed: int | None = None) -> str:
"""
Generate random data for this span and return its checksum.
Args:
vm: Virtual machine to run randstream on
dev: Device path (e.g., '/dev/xvdb')
seed: Optional seed for deterministic generation. If not provided,
caller must ensure seed is passed explicitly.
Returns:
Checksum of generated data as string
"""
seed_str = f'--seed {seed}' if seed is not None else ''
self.checksum = randstream(
vm, f'generate {seed_str} --position {self.position} --size {self.size} {dev}'.strip()
)
return self.checksum
def validate(self, vm: VM, dev: str) -> None:
"""
Validate random data for this span.
Args:
vm: Virtual machine to run randstream on
dev: Device path (e.g., '/dev/xvdb')
If checksum is set, validates against the expected checksum.
Otherwise, the stream itself contains checksums for each chunk
and will be validated using those internal checksums.
"""
expected_flags = f'--expected-checksum {self.checksum}' if self.checksum is not None else ''
randstream(vm, f'validate {expected_flags} --position {self.position} --size {self.size} {dev}')
def compute_span_layout(dev_size: int, total_size: int, num_spans: int, block_size: int) -> list[tuple[int, int]]:
"""
Compute the positions and sizes of spans across a device.
Spans are distributed such that the first span starts at 0 and the last
span ends at dev_size. Positions are computed right-to-left: each span's
start is rounded UP to the nearest multiple of align so the span sits
entirely within its slot. Bytes skipped by rounding are carried leftward
to the previous span, which absorbs them in its budget.
When total_size equals dev_size, spans are contiguous and cover the full
device with no gaps. When total_size is less than dev_size, spans are
evenly spread across the device with gaps between them.
Args:
dev_size: Total device size in bytes
total_size: Total number of bytes to write across all spans.
Use dev_size for full device coverage, or a smaller value
to leave gaps between spans.
num_spans: Number of spans to create
block_size: Block size in bytes to align span starts to.
All span starts are rounded up to the nearest multiple of align.
Bytes skipped by the rounding are carried to the previous span.
Returns:
List of (position, size) tuples, one per span.
Spans are guaranteed to:
- Not overlap
- Have first span starting at position 0
- Have last span ending at dev_size
- Sum of span sizes equals total_size
"""
per_span = total_size // num_spans
result: list[tuple[int, int]] = []
# Seed the remainder into the last span so it is distributed via carry
carry = total_size % num_spans
end = dev_size
for i in range(num_spans - 1, -1, -1):
budget = per_span + carry
if i == 0:
# First span always starts at 0; any leftover space before the
# next span becomes a gap (partial coverage only)
start = 0
size = min(budget, end)
else:
# Round start UP: this may shrink the span relative to budget;
# the trimmed bytes are carried left to the previous span
ideal_start = end - budget
start = ((ideal_start + block_size - 1) // block_size) * block_size
size = end - start
carry = budget - size
result.append((start, size))
end = start
result.reverse()
assert sum(s for _, s in result) == total_size
return result
def partially_populate_device(vm: VM, dev_path: str, dev_size: int, num_spans: int = 3, skip_spans: list[int] = []) \
-> list[StreamSpan]:
"""
Generate random data in multiple spans across a device.
Creates num_spans spans of random data distributed across the device.
Spans are positioned such that first span starts at 0 and last span
ends at dev_size, with middle spans evenly distributed in between.
ASCII visualization of span distribution:
For num_spans=3:
Device: [======================================]
Span 0: [****]
Span 1: [****]
Span 2: [****]
Gap: gap1 gap2 gap3
For num_spans=4 with skip_spans=[1]:
Device: [================================================]
Span 0: [**]
Span 1: [ ]
Span 2: [**]
Span 3: [**]
Args:
vm: Virtual machine to run randstream on
dev_path: Device path (e.g., '/dev/xvdb')
dev_size: Total device size in bytes
num_spans: Number of spans to create (default: 3)
skip_spans: List of span indices to skip (no data generated).
Skipped spans still exist in returned list with checksum=None.
(default: [])
Span alignment is controlled by the --write-volume-align pytest option.
Returns:
List of StreamSpan objects representing the generated spans.
"""
logging.info(f"Generate {dev_path} content")
total_size = min(dev_size, config.write_volume_cap)
# Validate skip_spans
assert all(0 <= i < num_spans for i in skip_spans), \
f"Invalid span index in skip_spans: must be 0 <= i < {num_spans}"
layout = compute_span_layout(dev_size, total_size, num_spans, config.write_volume_align)
spans: list[StreamSpan] = []
for i, (position, size) in enumerate(layout):
span = StreamSpan(position=position, size=size)
if i not in skip_spans:
span.generate(vm, dev_path, seed=1000 + i)
spans.append(span)
return spans
def validate_partially_populated_device(vm: VM, dev: str, spans: list[StreamSpan]) -> None:
logging.info(f"Validate {dev} content")
for span in spans:
if span.checksum is not None:
span.validate(vm, dev)