Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 18 additions & 3 deletions src/benchmark_radar/data_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from __future__ import annotations

import hashlib
import http.client
import ipaddress
import json
import os
Expand Down Expand Up @@ -250,12 +251,12 @@ def _manifest(
**({"If-None-Match": previous_etag} if previous_etag else {}),
}
request = urllib.request.Request(manifest_url, headers=headers)
response = self._open(request, allow_not_modified=True)
response = self._open(request, allow_not_modified=previous_etag is not None)
if response is None:
return None, previous_etag
with response:
try:
value = json.loads(response.read())
value = json.loads(self._read_body(response))
except json.JSONDecodeError as error:
raise DataSyncError(
"remote manifest is invalid JSON", code="invalid_manifest"
Expand Down Expand Up @@ -302,6 +303,20 @@ def _manifest(
raise DataSyncError("remote artifact file_count is invalid", code="invalid_manifest")
return value, etag

@staticmethod
def _read_body(response: Any, *args: int) -> bytes:
# The connection can fail after _open returns headers. Translate only
# read failures, keeping local disk-write errors distinct and avoiding
# partial response contents in the public error message.
try:
return response.read(*args)
except (OSError, http.client.HTTPException) as error:
raise DataSyncError(
f"remote transfer failed: {type(error).__name__}",
code="remote_unavailable",
status=503,
) from error

def _download(self, manifest: dict[str, Any], target: Path) -> None:
artifact = manifest["artifact"]
request = urllib.request.Request(
Expand All @@ -314,7 +329,7 @@ def _download(self, manifest: dict[str, Any], target: Path) -> None:
digest = hashlib.sha256()
size = 0
with self._open(request) as response, target.open("wb") as handle:
while chunk := response.read(1024 * 1024):
while chunk := self._read_body(response, 1024 * 1024):
size += len(chunk)
if size > artifact["size"]:
raise DataSyncError("download exceeds manifest size", code="invalid_artifact")
Expand Down
52 changes: 52 additions & 0 deletions tests/test_data_sync.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
from __future__ import annotations

import hashlib
import http.client
import io
import json
import os
Expand Down Expand Up @@ -275,6 +276,57 @@ def not_modified(request, **kwargs):
assert result["downloaded"] is False


@pytest.mark.parametrize("phase", ["manifest", "artifact"])
@pytest.mark.parametrize(
"failure", [http.client.IncompleteRead(b"partial"), ConnectionResetError()]
)
def test_sync_body_failure_is_structured_and_keeps_active_data(tmp_path, phase, failure):
# Headers can arrive before a transfer fails; the open-only exception
# handler leaked tracebacks instead of the CLI's machine-readable error.
first, bundle, url = _release(tmp_path / "old", day=29)
store = DataStore(
root=tmp_path / "home", manifest_url=url, urlopen=_Remote(url, first, bundle).urlopen
)
store.initialize()
previous = store.state_path.read_bytes()
next_manifest, next_bundle, _ = _release(tmp_path / "new", day=30)
remote = _Remote(url, next_manifest, next_bundle)

class InterruptedResponse(_Response):
def read(self, *args):
raise failure

def interrupted(request, **kwargs):
is_manifest = request.full_url == url
if is_manifest == (phase == "manifest"):
return InterruptedResponse(b"", url=request.full_url)
return remote.urlopen(request, **kwargs)

store.urlopen = interrupted
with pytest.raises(DataSyncError) as error:
store.sync()
assert error.value.code == "remote_unavailable"
assert error.value.status == 503
assert store.state_path.read_bytes() == previous
assert QueryService(store.query_paths()).status()["status"] == "ok"
assert not (store.root / ".download.tmp").exists()
assert not (store.root / "sync.lock").exists()


def test_init_rejects_unsolicited_not_modified_as_structured_error(tmp_path):
# There is no local release to reuse on a first install, so an unsolicited
# 304 must not flow into _current_result(None).
def not_modified(request, **kwargs):
assert "If-none-match" not in request.headers
raise urllib.error.HTTPError(request.full_url, 304, "Not Modified", {}, None)

store = DataStore(root=tmp_path / "home", urlopen=not_modified)
with pytest.raises(DataSyncError) as error:
store.initialize()
assert error.value.code == "remote_unavailable"
assert not store.state_path.exists()


def test_sync_does_not_call_current_when_active_data_is_corrupt(tmp_path: Path) -> None:
# Regression: manifest equality must not bless damaged local files as current.
manifest, bundle, manifest_url = _release(tmp_path)
Expand Down
Loading