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]>
154 lines
5.6 KiB
Python
154 lines
5.6 KiB
Python
"""Shared fixtures for the clawhdf5 Python binding tests."""
|
|
|
|
import os
|
|
import re
|
|
import threading
|
|
import time
|
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
|
|
|
import pytest
|
|
|
|
|
|
@pytest.fixture(scope="session")
|
|
def h5py():
|
|
"""h5py, or a skip — unless CLAWHDF5_REQUIRE_INTEROP=1, which makes a
|
|
missing h5py a failure (as for the Rust interop suites)."""
|
|
try:
|
|
import h5py as mod
|
|
except ImportError:
|
|
if os.environ.get("CLAWHDF5_REQUIRE_INTEROP") == "1":
|
|
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()
|