Allow retrieving compressed chunks

VDDK always decompresses the chunks it retrieves. However, some
callers may be interested in the compressed chunks, which may
be forwarded to another service involved in the backup process.

We'll add a read flag (skip_decompression). If set, the read
operation will return a "ReadResult" container, containing
a list of fragments and their compressed / decompressed lengths.

If the returned buffer size matches the AIO buffer size (which
now becomes configurable), at most one fragment will be returned.

While at it, we're adding perf tests that check various AIO buffer
sizes, cross checking against VDDK.
This commit is contained in:
Lucian Petrut
2026-09-17 15:29:11 +00:00
parent f0db6e166b
commit 2099591245
12 changed files with 914 additions and 94 deletions
+52 -3
View File
@@ -10,6 +10,9 @@ import pytest
from openvixdisklib import nfc_open
from tests.integration.base import SECTOR_SIZE, LabEnv, pattern_bytes
_1MIB = 1024 * 1024
_2MIB = 2 * 1024 * 1024
_16MIB = 16 * 1024 * 1024
_32MIB = 32 * 1024 * 1024
@@ -70,14 +73,60 @@ class TestNfcReadWrite:
[nfc_open.NFC_COMPRESSION_NONE, nfc_open.NFC_COMPRESSION_FASTLZ],
ids=["plain", "fastlz"],
)
def test_write_and_read_32mb(self, lab: LabEnv, compression: int) -> None:
"""Write 32 MiB (512 AIO chunks) and read it back in one request."""
@pytest.mark.parametrize(
"aio_buffer_count, aio_buffer_size",
[
(1, nfc_open.NFC_AIO_BUFFER_SIZE),
(1, _1MIB),
(1, _2MIB),
(4, _2MIB),
pytest.param(
1,
_16MIB,
marks=pytest.mark.xfail(
raises=nfc_open.NfcProtocolError,
reason="ESXi 8 rejects OPEN_SESSION bufSize 16 MiB",
strict=True,
),
),
pytest.param(
1,
_32MIB,
marks=pytest.mark.xfail(
raises=nfc_open.NfcProtocolError,
reason="ESXi 8 rejects OPEN_SESSION bufSize 32 MiB",
strict=True,
),
),
],
ids=[
"count1-64kib",
"count1-1mib",
"count1-2mib",
"count4-2mib",
"count1-16mib",
"count1-32mib",
],
)
def test_write_and_read_32mb(
self,
lab: LabEnv,
compression: int,
aio_buffer_count: int,
aio_buffer_size: int,
) -> None:
"""Write 32 MiB and read it back for several OPEN_SESSION sizes."""
n_sectors = _32MIB // SECTOR_SIZE
to_write = os.urandom(_32MIB)
with (
lab.authenticate(read_only=False) as session,
nfc_open.open_disk(
session, lab.disk_path, read_only=False, compression=compression
session,
lab.disk_path,
read_only=False,
compression=compression,
aio_buffer_size=aio_buffer_size,
aio_buffer_count=aio_buffer_count,
) as disk,
):
disk.write(0, n_sectors, to_write)
+86
View File
@@ -3,11 +3,15 @@
"""Exercise the VDDK-compatible openvixdisklib handle against the lab."""
from typing import Any
import pytest
from pyVim.connect import Disconnect
from pyVmomi import vim
from openvixdisklib import fastlz, nfc_open
from openvixdisklib import openvixdisklib as vixdisklib
from openvixdisklib.openvixdisklib import ReadResult
from tests.integration.base import (
SECTOR_AT_1GB,
SECTOR_SIZE,
@@ -28,6 +32,29 @@ def _virtual_disk_backing(
raise AssertionError(f"{vm._moId} has no virtual disk")
_2MIB = 2 * 1024 * 1024
def _rebuild_skip(buf: Any, result: ReadResult) -> bytes:
"""Decompress packed skip-decompression extras into uncompressed bytes."""
view = buf.raw if hasattr(buf, "raw") else buf
out = bytearray(result.uncompressed_length)
packed = 0
for frag in result.fragments:
extra = bytes(view[frag.offset : frag.offset + frag.length])
packed += frag.length
if frag.compression_type == nfc_open.NFC_COMPRESSION_FASTLZ:
chunk = fastlz.decompress(extra, frag.uncompressed_length)
elif frag.compression_type == nfc_open.NFC_COMPRESSION_NONE:
chunk = extra
else:
raise AssertionError(f"unexpected compression_type {frag.compression_type}")
assert len(chunk) == frag.uncompressed_length
out[frag.dest : frag.dest + frag.uncompressed_length] = chunk
assert packed == result.compressed_length
return bytes(out)
class TestOpenvixdisklib:
@pytest.mark.parametrize("transport_mode", ["nbdssl", "nbd"])
@pytest.mark.parametrize(
@@ -134,3 +161,62 @@ class TestOpenvixdisklib:
_wait_for_task(vm.RemoveAllSnapshots_Task())
finally:
Disconnect(si)
@pytest.mark.parametrize(
"aio_buffer_size, n_sectors, n_fragments",
[
(nfc_open.NFC_AIO_BUFFER_SIZE, 128, 1),
(nfc_open.NFC_AIO_BUFFER_SIZE, 129, 2),
(_2MIB, 129, 1),
],
ids=["64kib-128s", "64kib-129s", "2mib-129s"],
)
def test_skip_decompression_fastlz(
self,
lab: LabEnv,
aio_buffer_size: int,
n_sectors: int,
n_fragments: int,
) -> None:
"""Pack FastLZ extras and rebuild the same bytes as a normal read."""
length = n_sectors * SECTOR_SIZE
expected = pattern_bytes(length, b"OVDL-SKIP-")
handle = vixdisklib.VixDiskLibHandle(
vixdisklib_compatibility_version="8.0", config_path=None
)
write_buf = vixdisklib.get_buffer(length)
plain_buf = vixdisklib.get_buffer(length)
skip_buf = vixdisklib.get_buffer(length)
write_buf[:length] = expected
connect_kwargs = lab.vixdisklib_connect_kwargs(
{"allow_untrusted": lab.allow_untrusted, "transport_modes": "nbd"}
)
flags = vixdisklib.VIXDISKLIB_FLAG_OPEN_COMPRESSION_FASTLZ
with (
handle.connect(**connect_kwargs) as conn,
handle.open(
conn,
lab.disk_path,
flags=flags,
aio_buffer_size=aio_buffer_size,
aio_buffer_count=1,
) as disk,
):
handle.write(disk, 0, n_sectors, write_buf)
plain = handle.read(disk, 0, n_sectors, plain_buf)
skip = handle.read(disk, 0, n_sectors, skip_buf, skip_decompression=True)
assert isinstance(plain, ReadResult)
assert plain.fragments == ()
assert plain.uncompressed_length == length
assert plain.compressed_length <= length
assert skip.uncompressed_length == length
assert skip.compressed_length <= length
assert skip.compressed_length == plain.compressed_length
assert len(skip.fragments) == n_fragments
dests = {frag.dest for frag in skip.fragments}
if n_fragments == 1:
assert dests == {0}
else:
assert dests == {0, nfc_open.NFC_AIO_BUFFER_SIZE}
assert plain_buf.raw[:length] == expected
assert _rebuild_skip(skip_buf, skip) == expected
+266 -34
View File
@@ -5,18 +5,71 @@
from __future__ import annotations
import json
import os
import pickle
import subprocess
import sys
import tempfile
import time
from typing import Any
from openvixdisklib import nfc_open
from openvixdisklib import openvixdisklib as open_vix
from tests.integration import vixdisklib
from tests.integration.base import SECTOR_SIZE, LabEnv, pattern_bytes
from tests.integration.base import (
SECTOR_SIZE,
LabEnv,
ensure_vddk_library_path,
pattern_bytes,
)
_REPO = os.path.abspath(os.path.join(os.path.dirname(__file__), "../.."))
_SIZES = (
("64KiB", 64 * 1024),
("129 sectors", 129 * SECTOR_SIZE),
("32MiB", 32 * 1024 * 1024),
)
_1MIB = 1024 * 1024
_2MIB = 2 * 1024 * 1024
# ESXi 8 accepts 64 KiB, 1 MiB, and 2 MiB extras. 16 MiB and 32 MiB
# OPEN_SESSION are rejected (AIO error). Broadcom's 16 MiB cap is
# size×count session memory, not a larger extra; VDDK's per-buffer max
# is 2 MiB (``BufSizeIn64KB=16`` is 1 MiB).
_AIO_SESSIONS = (
(nfc_open.NFC_AIO_BUFFER_COUNT, nfc_open.NFC_AIO_BUFFER_SIZE),
(1, _1MIB),
(1, _2MIB),
(4, _2MIB),
)
_64KIB = 64 * 1024
def _aio_size_label(nbytes: int) -> str:
"""Return a short label for an AIO extra size."""
if nbytes % (1024 * 1024) == 0:
return f"{nbytes // (1024 * 1024)}MiB"
if nbytes % 1024 == 0:
return f"{nbytes // 1024}KiB"
return str(nbytes)
def _vddk_aio_config(
directory: str, aio_buffer_size: int, aio_buffer_count: int
) -> str:
"""Write a temp VDDK config for ``BufSizeIn64KB`` and ``BufCount``."""
if aio_buffer_size % _64KIB:
raise ValueError(
f"VDDK BufSizeIn64KB needs a 64 KiB multiple, got {aio_buffer_size}"
)
path = os.path.join(directory, "vddk.config")
with open(path, "w", encoding="utf-8") as config:
config.write(f"tmpDirectory={directory}\n")
config.write(
f"vixDiskLib.nfcAio.Session.BufSizeIn64KB={aio_buffer_size // _64KIB}\n"
)
config.write(f"vixDiskLib.nfcAio.Session.BufCount={aio_buffer_count}\n")
return path
def _connect_extra(lab: LabEnv, module: Any, transport_mode: str) -> dict[str, Any]:
@@ -27,37 +80,182 @@ def _connect_extra(lab: LabEnv, module: Any, transport_mode: str) -> dict[str, A
return extra
def _time_write_read(
def _time_write_read_once(
lab: LabEnv,
module: Any,
payload: bytes,
flags: int = 0,
transport_mode: str = "nbdssl",
aio_buffer_size: int = nfc_open.NFC_AIO_BUFFER_SIZE,
aio_buffer_count: int = nfc_open.NFC_AIO_BUFFER_COUNT,
config_dir: str | None = None,
skip_decompression: bool = False,
) -> tuple[float, float]:
"""Write ``payload`` at sector 0, read it back, and return durations."""
"""Write ``payload`` at sector 0, read it back, and return durations.
Does not call ``VixDiskLib_Exit``. Native VDDK double-frees if
``InitEx``/``Exit`` are paired more than once in the same process.
``skip_decompression`` is OpenVixDiskLib FastLZ skip; ``buf`` then
holds packed extras, not sector bytes.
"""
n_sectors = len(payload) // SECTOR_SIZE
config_path = None
if module is vixdisklib:
if config_dir is None:
raise ValueError("VDDK timings need a config_dir")
config_path = _vddk_aio_config(config_dir, aio_buffer_size, aio_buffer_count)
handle = module.VixDiskLibHandle(
vixdisklib_compatibility_version="8.0", config_path=None
vixdisklib_compatibility_version="8.0", config_path=config_path
)
write_buf = module.get_buffer(len(payload))
read_buf = module.get_buffer(len(payload))
write_buf[: len(payload)] = payload
kwargs = lab.vixdisklib_connect_kwargs(_connect_extra(lab, module, transport_mode))
open_kwargs: dict[str, Any] = {"flags": flags}
if module is open_vix:
open_kwargs["aio_buffer_size"] = aio_buffer_size
open_kwargs["aio_buffer_count"] = aio_buffer_count
with (
handle.connect(**kwargs) as conn,
handle.open(conn, lab.disk_path, flags=flags) as disk,
handle.open(conn, lab.disk_path, **open_kwargs) as disk,
):
started = time.perf_counter()
handle.write(disk, 0, n_sectors, write_buf)
write_s = time.perf_counter() - started
read_buf[: len(payload)] = b"\xa5" * len(payload)
read_kwargs: dict[str, Any] = {}
if skip_decompression:
read_kwargs["skip_decompression"] = True
started = time.perf_counter()
handle.read(disk, 0, n_sectors, read_buf)
result = handle.read(disk, 0, n_sectors, read_buf, **read_kwargs)
read_s = time.perf_counter() - started
assert read_buf.raw[: len(payload)] == payload
if skip_decompression:
assert result.uncompressed_length == len(payload)
assert result.compressed_length <= len(payload)
assert result.fragments
else:
assert read_buf.raw[: len(payload)] == payload
return write_s, read_s
def _run_vddk_worker(lab_pkl: str, job_pkl: str, work_dir: str) -> None:
"""InitEx once in this process, time one write/read, write result.json."""
os.environ.pop("LD_PRELOAD", None)
ensure_vddk_library_path()
with open(lab_pkl, "rb") as pickle_file:
lab = pickle.load(pickle_file)
with open(job_pkl, "rb") as pickle_file:
job = pickle.load(pickle_file)
payload = pattern_bytes(job["nbytes"], f"PERF-{job['label']}-".encode())
write_s, read_s = _time_write_read_once(
lab,
vixdisklib,
payload,
flags=job["flags"],
transport_mode=job["transport_mode"],
aio_buffer_size=job["aio_buffer_size"],
aio_buffer_count=job["aio_buffer_count"],
config_dir=work_dir,
)
result_path = os.path.join(work_dir, "result.json")
with open(result_path, "w", encoding="utf-8") as result_file:
json.dump({"write_s": write_s, "read_s": read_s}, result_file)
def _time_vddk_subprocess(
lab: LabEnv,
label: str,
nbytes: int,
flags: int,
transport_mode: str,
aio_buffer_size: int,
aio_buffer_count: int,
) -> tuple[float, float]:
"""Time native VDDK in a child process so InitEx sees this AIO config."""
with tempfile.TemporaryDirectory(prefix="vddk-perf-") as work_dir:
lab_pkl = os.path.join(work_dir, "lab.pkl")
job_pkl = os.path.join(work_dir, "job.pkl")
result_path = os.path.join(work_dir, "result.json")
with open(lab_pkl, "wb") as pickle_file:
pickle.dump(lab, pickle_file)
with open(job_pkl, "wb") as pickle_file:
pickle.dump(
{
"label": label,
"nbytes": nbytes,
"flags": flags,
"transport_mode": transport_mode,
"aio_buffer_size": aio_buffer_size,
"aio_buffer_count": aio_buffer_count,
},
pickle_file,
)
env = os.environ.copy()
env.pop("LD_PRELOAD", None)
pythonpath = env.get("PYTHONPATH", "")
env["PYTHONPATH"] = _REPO if not pythonpath else f"{_REPO}:{pythonpath}"
proc = subprocess.run(
[
sys.executable,
os.path.abspath(__file__),
"--vddk-worker",
lab_pkl,
job_pkl,
work_dir,
],
check=False,
capture_output=True,
text=True,
env=env,
cwd=_REPO,
)
if proc.returncode != 0 or not os.path.exists(result_path):
raise RuntimeError(
"VDDK perf worker failed "
f"(exit {proc.returncode}): {proc.stderr}\n{proc.stdout}"
)
with open(result_path, encoding="utf-8") as result_file:
result = json.load(result_file)
return float(result["write_s"]), float(result["read_s"])
def _time_write_read(
lab: LabEnv,
module: Any,
label: str,
nbytes: int,
flags: int = 0,
transport_mode: str = "nbdssl",
aio_buffer_size: int = nfc_open.NFC_AIO_BUFFER_SIZE,
aio_buffer_count: int = nfc_open.NFC_AIO_BUFFER_COUNT,
skip_decompression: bool = False,
) -> tuple[float, float]:
"""Time one write/read; native VDDK runs in a subprocess."""
if skip_decompression and module is vixdisklib:
raise ValueError("skip_decompression is OpenVixDiskLib-only")
if module is vixdisklib:
return _time_vddk_subprocess(
lab,
label,
nbytes,
flags,
transport_mode,
aio_buffer_size,
aio_buffer_count,
)
payload = pattern_bytes(nbytes, f"PERF-{label}-".encode())
return _time_write_read_once(
lab,
module,
payload,
flags=flags,
transport_mode=transport_mode,
aio_buffer_size=aio_buffer_size,
aio_buffer_count=aio_buffer_count,
skip_decompression=skip_decompression,
)
def _mib_per_s(nbytes: int, seconds: float) -> float:
if seconds <= 0:
return float("inf")
@@ -66,49 +264,74 @@ def _mib_per_s(nbytes: int, seconds: float) -> float:
class TestCompare:
def test_write_read_throughput(self, lab: LabEnv, vddk: None) -> None:
"""Time matching write/read sizes on VDDK and openvixdisklib."""
"""Time matching write/read sizes on VDDK and openvixdisklib.
Prints ``aio_size`` / ``aio_count`` for each OPEN_SESSION
(64 KiB×1, 1 MiB×1, 2 MiB×1, 2 MiB×4). VDDK gets those via
``vixDiskLib.nfcAio.Session.BufSizeIn64KB`` / ``BufCount`` in a
fresh process per row (``VixDiskLib_Exit`` is not loop-safe).
``fastlz-skip`` is OpenVixDiskLib ``skip_decompression`` (packed
extras, no FastLZ decode); VDDK has no equivalent.
"""
libraries = (
("vddk", vixdisklib),
("openvixdisklib", open_vix),
)
transports = ("nbdssl", "nbd")
open_modes = (
("plain", 0),
("fastlz", vixdisklib.VIXDISKLIB_FLAG_OPEN_COMPRESSION_FASTLZ),
("plain", 0, False),
("fastlz", vixdisklib.VIXDISKLIB_FLAG_OPEN_COMPRESSION_FASTLZ, False),
(
"fastlz-skip",
vixdisklib.VIXDISKLIB_FLAG_OPEN_COMPRESSION_FASTLZ,
True,
),
)
rows: list[tuple[str, str, str, str, float, float, float, float]] = []
rows: list[tuple[str, str, int, str, str, str, float, float, float, float]] = []
for label, nbytes in _SIZES:
payload = pattern_bytes(nbytes, f"PERF-{label}-".encode())
for transport_mode in transports:
for mode_name, flags in open_modes:
for name, module in libraries:
write_s, read_s = _time_write_read(
lab,
module,
payload,
flags=flags,
transport_mode=transport_mode,
)
rows.append(
(
for aio_buffer_count, aio_buffer_size in _AIO_SESSIONS:
aio_label = _aio_size_label(aio_buffer_size)
for transport_mode in transports:
for mode_name, flags, skip_decompression in open_modes:
for name, module in libraries:
if skip_decompression and module is vixdisklib:
continue
write_s, read_s = _time_write_read(
lab,
module,
label,
transport_mode,
mode_name,
name,
write_s,
read_s,
_mib_per_s(nbytes, write_s),
_mib_per_s(nbytes, read_s),
nbytes,
flags=flags,
transport_mode=transport_mode,
aio_buffer_size=aio_buffer_size,
aio_buffer_count=aio_buffer_count,
skip_decompression=skip_decompression,
)
rows.append(
(
label,
aio_label,
aio_buffer_count,
transport_mode,
mode_name,
name,
write_s,
read_s,
_mib_per_s(nbytes, write_s),
_mib_per_s(nbytes, read_s),
)
)
)
print()
print(
f"{'size':<14} {'transport':<10} {'flags':<8} {'library':<16} "
f"{'size':<14} {'aio_size':<8} {'aio_count':>9} "
f"{'transport':<10} {'flags':<12} {'library':<16} "
f"{'write_s':>10} {'read_s':>10} "
f"{'write_MiB/s':>12} {'read_MiB/s':>12}"
)
for (
label,
aio_label,
aio_buffer_count,
transport_mode,
mode_name,
name,
@@ -118,7 +341,16 @@ class TestCompare:
read_r,
) in rows:
print(
f"{label:<14} {transport_mode:<10} {mode_name:<8} {name:<16} "
f"{label:<14} {aio_label:<8} {aio_buffer_count:>9} "
f"{transport_mode:<10} {mode_name:<12} {name:<16} "
f"{write_s:10.3f} {read_s:10.3f} "
f"{write_r:12.1f} {read_r:12.1f}"
)
if __name__ == "__main__":
if sys.argv[1:2] == ["--vddk-worker"]:
_, _, lab_pkl, job_pkl, work_dir = sys.argv
_run_vddk_worker(lab_pkl, job_pkl, work_dir)
else:
raise SystemExit("usage: test_compare.py --vddk-worker LAB JOB DIR")