py: remote files (clawhdf5.File(url), File.open_url) through File::storage()
The Python bindings could not open a remote file: they parsed through File::as_bytes() in eight places (path lookups, object headers, dataspaces, attributes, group listings, the global heap of variable-length data), which a storage-backed file does not have. - Every object of a File now shares one handle (src/handle.rs) that runs all file access, metadata included, with the GIL released and parses through File::storage() and the clawhdf5_format *_in functions. Local files take the same path (their storage is the mmap). - clawhdf5.File(url) opens any scheme://... through clawhdf5_remote::storage_for_url (read-only; another mode is a ValueError). File.open_url(url, **options) takes the cache and HTTP options (block_size, cache_size, headers, retries, timeout, allow_full_download, max_full_download, require_validator, max_redirects, max_parallel); File.remote_stats gives the block cache's counters. - Default build: plain HTTP only, no C. https (rustls/ring) and s3/gcs/azure (aws-lc-rs) are opt-in features of clawhdf5-py, and ci-test.sh's no-C check now covers the crate. - A failed storage read (network error, file changed on the server) is an OSError, never KeyError/ValueError and never data; `key in group` raises it instead of answering False. Tests: the read-vs-h5py suite runs locally and over HTTP (1 MiB and 1 KiB blocks) against a range-capable http.server in the test process (conftest.RangeServer); test_remote.py covers request counts, cache hits, a server without Range support, a changed file, a server that hangs up, 16 threads, and a spinning thread that keeps running while a read waits on 0.2 s requests. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -1,6 +1,10 @@
|
||||
"""Shared fixtures for the clawhdf5 Python binding tests."""
|
||||
|
||||
import os
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||
|
||||
import pytest
|
||||
|
||||
@@ -16,3 +20,134 @@ def h5py():
|
||||
pytest.fail("h5py is required (CLAWHDF5_REQUIRE_INTEROP=1) but not importable")
|
||||
pytest.skip("h5py not installed")
|
||||
return mod
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# An HTTP server for remote reads
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
_RANGE = re.compile(r"^bytes=(\d*)-(\d*)$")
|
||||
|
||||
|
||||
class RangeServer:
|
||||
"""A static file server on 127.0.0.1, in a thread of this process, that
|
||||
answers `Range: bytes=a-b` with 206 and `Content-Range` (the way S3 and
|
||||
common web servers do), sends an ETag and honours `If-Match`.
|
||||
|
||||
- `ranges=False`: ignores `Range` and answers 200 with the whole file,
|
||||
like a server without range support.
|
||||
- `down` (set by `close()`): hang up on every request.
|
||||
- `delay`: seconds to wait before answering each request after the
|
||||
first `delay_after` ones (a slow network).
|
||||
- `log`: every request as `(method, path, range header)`.
|
||||
"""
|
||||
|
||||
def __init__(self, root, ranges=True):
|
||||
self.root = str(root)
|
||||
self.ranges = ranges
|
||||
self.delay = 0.0
|
||||
self.delay_after = 0
|
||||
self.down = False
|
||||
self.log = []
|
||||
self._lock = threading.Lock()
|
||||
server = self
|
||||
|
||||
class Handler(BaseHTTPRequestHandler):
|
||||
protocol_version = "HTTP/1.1"
|
||||
|
||||
def log_message(self, *args): # quiet
|
||||
pass
|
||||
|
||||
def do_HEAD(self):
|
||||
self._serve(body=False)
|
||||
|
||||
def do_GET(self):
|
||||
self._serve(body=True)
|
||||
|
||||
def _serve(self, body):
|
||||
with server._lock:
|
||||
server.log.append((self.command, self.path, self.headers.get("Range")))
|
||||
n = len(server.log)
|
||||
if server.delay and n > server.delay_after:
|
||||
time.sleep(server.delay)
|
||||
if server.down:
|
||||
# Hang up without an answer (open keep-alive
|
||||
# connections outlive shutdown(), so close() sets this).
|
||||
self.close_connection = True
|
||||
return
|
||||
path = os.path.join(server.root, self.path.lstrip("/").split("?")[0])
|
||||
if not os.path.isfile(path):
|
||||
self.send_response(404)
|
||||
self.send_header("Content-Length", "0")
|
||||
self.end_headers()
|
||||
return
|
||||
with open(path, "rb") as fh:
|
||||
data = fh.read()
|
||||
st = os.stat(path)
|
||||
etag = f'"{st.st_mtime_ns:x}-{st.st_size:x}"'
|
||||
want = self.headers.get("If-Match")
|
||||
if want is not None and want != etag and want != "*":
|
||||
self.send_response(412)
|
||||
self.send_header("Content-Length", "0")
|
||||
self.end_headers()
|
||||
return
|
||||
rng = self.headers.get("Range") if server.ranges else None
|
||||
m = _RANGE.match(rng.strip()) if rng else None
|
||||
if m and (m.group(1) or m.group(2)):
|
||||
size = len(data)
|
||||
if m.group(1):
|
||||
start = int(m.group(1))
|
||||
end = int(m.group(2)) if m.group(2) else size - 1
|
||||
else:
|
||||
start = max(0, size - int(m.group(2)))
|
||||
end = size - 1
|
||||
if start >= size:
|
||||
self.send_response(416)
|
||||
self.send_header("Content-Range", f"bytes */{size}")
|
||||
self.send_header("Content-Length", "0")
|
||||
self.end_headers()
|
||||
return
|
||||
end = min(end, size - 1)
|
||||
part = data[start : end + 1]
|
||||
self.send_response(206)
|
||||
self.send_header("Content-Range", f"bytes {start}-{end}/{size}")
|
||||
else:
|
||||
part = data
|
||||
self.send_response(200)
|
||||
if server.ranges:
|
||||
self.send_header("Accept-Ranges", "bytes")
|
||||
self.send_header("ETag", etag)
|
||||
self.send_header("Content-Length", str(len(part)))
|
||||
self.send_header("Content-Type", "application/x-hdf5")
|
||||
self.end_headers()
|
||||
if body:
|
||||
try:
|
||||
self.wfile.write(part)
|
||||
except (BrokenPipeError, ConnectionResetError):
|
||||
pass
|
||||
|
||||
self.httpd = ThreadingHTTPServer(("127.0.0.1", 0), Handler)
|
||||
self.httpd.daemon_threads = True
|
||||
self.port = self.httpd.server_address[1]
|
||||
self.thread = threading.Thread(target=self.httpd.serve_forever, daemon=True)
|
||||
self.thread.start()
|
||||
|
||||
def url(self, name):
|
||||
return f"http://127.0.0.1:{self.port}/{name}"
|
||||
|
||||
def requests(self):
|
||||
with self._lock:
|
||||
return len(self.log)
|
||||
|
||||
def close(self):
|
||||
self.down = True
|
||||
self.httpd.shutdown()
|
||||
self.httpd.server_close()
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def range_server(tmp_path):
|
||||
"""A range-capable server over `tmp_path`."""
|
||||
server = RangeServer(tmp_path)
|
||||
yield server
|
||||
server.close()
|
||||
|
||||
@@ -178,15 +178,32 @@ def _write_fixture(h5py, path):
|
||||
g.attrs["depth"] = np.int8(3)
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def pair(h5py, tmp_path_factory):
|
||||
path = str(tmp_path_factory.mktemp("h5") / "fixture.h5")
|
||||
@pytest.fixture(scope="module", params=["local", "http", "http-1k-blocks"])
|
||||
def pair(request, h5py, tmp_path_factory):
|
||||
"""The fixture file through h5py and through clawhdf5: opened locally,
|
||||
and over HTTP range requests (a local server in this process) with the
|
||||
default 1 MiB blocks and with 1 KiB blocks, so every structure is read
|
||||
through many small ranges."""
|
||||
from conftest import RangeServer
|
||||
|
||||
root = tmp_path_factory.mktemp("h5")
|
||||
path = str(root / "fixture.h5")
|
||||
_write_fixture(h5py, path)
|
||||
theirs = h5py.File(path, "r")
|
||||
ours = clawhdf5.File(path, "r")
|
||||
server = None
|
||||
if request.param == "local":
|
||||
ours = clawhdf5.File(path, "r")
|
||||
else:
|
||||
server = RangeServer(root)
|
||||
if request.param == "http":
|
||||
ours = clawhdf5.File(server.url("fixture.h5"))
|
||||
else:
|
||||
ours = clawhdf5.File.open_url(server.url("fixture.h5"), block_size=1024)
|
||||
yield ours, theirs, path
|
||||
theirs.close()
|
||||
ours.close()
|
||||
if server is not None:
|
||||
server.close()
|
||||
|
||||
|
||||
def _all_datasets(h5py, f):
|
||||
|
||||
@@ -0,0 +1,244 @@
|
||||
"""Remote files: `clawhdf5.File(url)` / `File.open_url(url, ...)` read over
|
||||
HTTP range requests (clawhdf5-remote's block cache), against a server in this
|
||||
process (conftest.RangeServer). Values are compared with h5py reading the
|
||||
same file locally; the rest checks what the server saw (only the blocks a
|
||||
read needs are fetched), the failure modes (no range support, a missing
|
||||
file, a file that changes, a server that goes away: errors, never wrong
|
||||
data), and that the GIL is released while a read waits on the network."""
|
||||
|
||||
import os
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
|
||||
import numpy as np
|
||||
import pytest
|
||||
|
||||
import clawhdf5
|
||||
from conftest import RangeServer
|
||||
|
||||
|
||||
def _write(h5py, path):
|
||||
rng = np.random.default_rng(7)
|
||||
with h5py.File(path, "w") as f:
|
||||
f.create_dataset("contig", data=rng.standard_normal((400, 300)))
|
||||
f.create_dataset(
|
||||
"chunked",
|
||||
data=rng.integers(0, 1000, size=(512, 512), dtype="<i4"),
|
||||
chunks=(64, 64),
|
||||
compression="gzip",
|
||||
)
|
||||
f.create_dataset("strings", data=["alpha", "beta", "gamma"], dtype=h5py.string_dtype())
|
||||
g = f.create_group("grp")
|
||||
g.create_dataset("small", data=np.arange(10, dtype="<u2"))
|
||||
g.attrs["units"] = "m/s"
|
||||
f.attrs["title"] = "remote test"
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def remote_file(h5py, tmp_path, range_server):
|
||||
path = tmp_path / "remote.h5"
|
||||
_write(h5py, str(path))
|
||||
return path, range_server.url("remote.h5"), range_server
|
||||
|
||||
|
||||
def test_remote_reads_match_h5py(h5py, remote_file):
|
||||
path, url, server = remote_file
|
||||
with h5py.File(path, "r") as theirs, clawhdf5.File(url) as ours:
|
||||
assert ours.filename == url
|
||||
assert list(ours.keys()) == list(theirs.keys())
|
||||
assert "grp/small" in ours and "nope" not in ours
|
||||
for name in ["contig", "chunked", "grp/small"]:
|
||||
np.testing.assert_array_equal(ours[name][...], theirs[name][...])
|
||||
np.testing.assert_array_equal(ours[name][3:7], theirs[name][3:7])
|
||||
np.testing.assert_array_equal(ours["chunked"][100:130, 200:300:3], theirs["chunked"][100:130, 200:300:3])
|
||||
np.testing.assert_array_equal(ours["chunked"][[1, 70, 300], 5], theirs["chunked"][[1, 70, 300], 5])
|
||||
assert list(ours["strings"][...]) == list(theirs["strings"][...])
|
||||
assert ours["grp"].attrs["units"] == theirs["grp"].attrs["units"]
|
||||
assert ours.attrs["title"] == theirs.attrs["title"]
|
||||
assert ours["chunked"].maxshape == theirs["chunked"].maxshape
|
||||
assert server.requests() >= 2
|
||||
|
||||
|
||||
def test_a_small_read_fetches_only_its_blocks(h5py, remote_file):
|
||||
"""With 4 KiB blocks, opening and reading one chunk of a 1 MB chunked
|
||||
dataset costs a handful of requests and a few blocks, not the file."""
|
||||
path, url, server = remote_file
|
||||
size = os.path.getsize(path)
|
||||
f = clawhdf5.File.open_url(url, block_size=4096)
|
||||
opened = server.requests()
|
||||
assert opened == 1, server.log
|
||||
ds = f["chunked"]
|
||||
got = ds[0:10, 0:10]
|
||||
with h5py.File(path, "r") as theirs:
|
||||
np.testing.assert_array_equal(got, theirs["chunked"][0:10, 0:10])
|
||||
stats = f.remote_stats
|
||||
assert stats["bytes_fetched"] < size / 4, (stats, size)
|
||||
assert server.requests() - opened <= 12, server.log
|
||||
# A second read of the same region is served by the cache.
|
||||
before = server.requests()
|
||||
ds[0:10, 0:10]
|
||||
assert server.requests() == before
|
||||
assert f.remote_stats["hits"] > stats["hits"]
|
||||
assert clawhdf5.File(str(path), "r").remote_stats is None
|
||||
|
||||
|
||||
def test_server_without_range_support(h5py, tmp_path):
|
||||
"""A server that ignores Range answers 200 with the whole file: that is
|
||||
an OSError by default, and a whole download when allowed."""
|
||||
path = tmp_path / "remote.h5"
|
||||
_write(h5py, str(path))
|
||||
server = RangeServer(tmp_path, ranges=False)
|
||||
try:
|
||||
url = server.url("remote.h5")
|
||||
with pytest.raises(OSError, match="range"):
|
||||
clawhdf5.File(url)
|
||||
with clawhdf5.File.open_url(url, allow_full_download=True) as ours, h5py.File(path, "r") as theirs:
|
||||
np.testing.assert_array_equal(ours["chunked"][...], theirs["chunked"][...])
|
||||
np.testing.assert_array_equal(ours["contig"][5], theirs["contig"][5])
|
||||
with pytest.raises(OSError):
|
||||
clawhdf5.File.open_url(url, allow_full_download=True, max_full_download=1000)
|
||||
finally:
|
||||
server.close()
|
||||
|
||||
|
||||
def test_errors_are_oserrors(remote_file):
|
||||
_, url, server = remote_file
|
||||
with pytest.raises(OSError, match="404"):
|
||||
clawhdf5.File(server.url("missing.h5"))
|
||||
with pytest.raises(ValueError, match="read-only"):
|
||||
clawhdf5.File(url, "r+")
|
||||
with pytest.raises(ValueError, match="read-only"):
|
||||
clawhdf5.File(url, "w")
|
||||
with pytest.raises(OSError, match="unsupported URL"):
|
||||
clawhdf5.File("nosuchscheme://x/y.h5")
|
||||
with pytest.raises(ValueError):
|
||||
clawhdf5.File.open_url(url, block_size=0)
|
||||
with pytest.raises(TypeError):
|
||||
clawhdf5.File.open_url(url, no_such_option=1)
|
||||
|
||||
|
||||
def test_object_store_urls_need_their_features():
|
||||
"""The default wheel has no S3/GCS/Azure clients (aws-lc-rs builds C):
|
||||
such a URL is an OSError naming the build feature."""
|
||||
for url, feature in [("s3://bucket/k.h5", "s3"), ("gs://b/k.h5", "gcs"), ("az://c/k.h5", "azure")]:
|
||||
try:
|
||||
clawhdf5.File(url)
|
||||
except OSError as e:
|
||||
if "feature" in str(e):
|
||||
assert f"`{feature}`" in str(e), str(e)
|
||||
else:
|
||||
pytest.fail(f"{url} opened")
|
||||
|
||||
|
||||
def test_https_needs_the_https_feature():
|
||||
"""The default wheel has no TLS stack (rustls needs ring, which builds C):
|
||||
an https URL is an OSError that names the build feature."""
|
||||
with pytest.raises(OSError) as e:
|
||||
clawhdf5.File("https://127.0.0.1:1/x.h5")
|
||||
msg = str(e.value)
|
||||
# Built with `--features https` the error is the refused connection.
|
||||
assert "https" in msg or "connect" in msg.lower() or "refused" in msg.lower(), msg
|
||||
|
||||
|
||||
def test_a_changed_file_is_an_error_not_mixed_data(h5py, remote_file):
|
||||
path, url, _ = remote_file
|
||||
f = clawhdf5.File.open_url(url, block_size=1024)
|
||||
first = f["grp/small"][...]
|
||||
# Rewrite the file with other values: new ETag, same name.
|
||||
time.sleep(0.01)
|
||||
with h5py.File(path, "w") as g:
|
||||
g.create_dataset("contig", data=np.zeros((400, 300)))
|
||||
with pytest.raises(OSError, match="changed"):
|
||||
f["contig"][...]
|
||||
np.testing.assert_array_equal(first, np.arange(10, dtype="<u2"))
|
||||
|
||||
|
||||
def test_a_server_that_goes_away_is_an_error(h5py, tmp_path):
|
||||
path = tmp_path / "remote.h5"
|
||||
_write(h5py, str(path))
|
||||
server = RangeServer(tmp_path)
|
||||
f = clawhdf5.File.open_url(server.url("remote.h5"), block_size=1024, retries=0, timeout=2)
|
||||
ds = f["contig"]
|
||||
server.close()
|
||||
with pytest.raises(OSError):
|
||||
ds[...]
|
||||
|
||||
|
||||
def test_threads_read_one_remote_file(h5py, remote_file):
|
||||
path, url, _ = remote_file
|
||||
f = clawhdf5.File.open_url(url, block_size=2048)
|
||||
with h5py.File(path, "r") as theirs:
|
||||
expected = theirs["chunked"][...]
|
||||
errors = []
|
||||
|
||||
def work(i):
|
||||
try:
|
||||
rows = slice((i * 37) % 400, (i * 37) % 400 + 64)
|
||||
np.testing.assert_array_equal(f["chunked"][rows], expected[rows])
|
||||
except Exception as e: # noqa: BLE001
|
||||
errors.append(e)
|
||||
|
||||
threads = [threading.Thread(target=work, args=(i,)) for i in range(16)]
|
||||
for t in threads:
|
||||
t.start()
|
||||
for t in threads:
|
||||
t.join()
|
||||
assert not errors, errors[:3]
|
||||
|
||||
|
||||
def test_remote_reads_release_the_gil(h5py, remote_file):
|
||||
"""A read waiting on a slow server lets other Python threads run: a
|
||||
thread counting in a loop keeps counting (and never stalls for long)
|
||||
while the main thread reads through requests that each take 0.2 s."""
|
||||
_, url, server = remote_file
|
||||
f = clawhdf5.File.open_url(url, block_size=1024, max_parallel=1)
|
||||
ds = f["contig"]
|
||||
server.delay = 0.2
|
||||
server.delay_after = server.requests()
|
||||
old = sys.getswitchinterval()
|
||||
sys.setswitchinterval(0.001)
|
||||
stop = threading.Event()
|
||||
progress = {"n": 0, "worst": 0.0}
|
||||
|
||||
def spin():
|
||||
last = time.perf_counter()
|
||||
while not stop.is_set():
|
||||
now = time.perf_counter()
|
||||
progress["worst"] = max(progress["worst"], now - last)
|
||||
last = now
|
||||
progress["n"] += 1
|
||||
|
||||
t = threading.Thread(target=spin)
|
||||
try:
|
||||
t.start()
|
||||
time.sleep(0.02)
|
||||
t0 = time.perf_counter()
|
||||
before = server.requests()
|
||||
ds[0:2]
|
||||
took = time.perf_counter() - t0
|
||||
stop.set()
|
||||
t.join()
|
||||
finally:
|
||||
sys.setswitchinterval(old)
|
||||
assert server.requests() > before
|
||||
assert took >= 0.2, took
|
||||
assert progress["n"] > 1000, progress
|
||||
# Held across a 0.2 s request, the spinner would stall that long.
|
||||
assert progress["worst"] < 0.1, (progress, took)
|
||||
|
||||
|
||||
def test_a_clawhdf5_written_file_reads_the_same_remotely(tmp_path, range_server):
|
||||
path = tmp_path / "ours.h5"
|
||||
data = np.arange(3000, dtype="<f8").reshape(100, 30)
|
||||
with clawhdf5.File(str(path), "w") as f:
|
||||
f.create_dataset("d", data=data, chunks=[10, 30], compression="gzip")
|
||||
g = f.create_group("g")
|
||||
g.create_dataset("i", data=np.arange(5, dtype="<i4"))
|
||||
f.attrs["k"] = 3
|
||||
with clawhdf5.File(range_server.url("ours.h5")) as f, clawhdf5.File(str(path)) as local:
|
||||
np.testing.assert_array_equal(f["d"][...], data)
|
||||
np.testing.assert_array_equal(f["d"][5:9, ::4], local["d"][5:9, ::4])
|
||||
np.testing.assert_array_equal(f["g/i"][...], np.arange(5))
|
||||
assert f.attrs["k"] == local.attrs["k"]
|
||||
assert "file (read" in repr(f).lower() and "127.0.0.1" in repr(f)
|
||||
Reference in New Issue
Block a user