Skip to content

Commit f977912

Browse files
committed
Fix exists() exception narrowing and add real HF revision integration test
FsspecInputFile.exists() reimplemented exists()-via-info() with a try/except catching only FileNotFoundError, unlike fsspec's own lexists()/exists() which swallow any exception. Delegate to self._fs.exists() directly instead, which already forwards kwargs to info() with the correct broad exception handling. Also add an HF_TOKEN-gated integration test against a real (temporary) Hugging Face dataset repo, demonstrating the concrete problem this property solves: without hf.revision, reads always follow the repo's moving default branch, so a file "written" at one commit silently returns different content once someone pushes a new commit to the same path. Minor: use a _HF_SCHEMES frozenset in _get_fs_kwargs for consistency with the existing _ADLS_SCHEMES dispatch pattern, and de-duplicate/strengthen two tests that asserted only on the private _fs_kwargs attribute.
1 parent 4057eb3 commit f977912

2 files changed

Lines changed: 53 additions & 12 deletions

File tree

‎pyiceberg/io/fsspec.py‎

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -320,6 +320,7 @@ def _hf(properties: Properties) -> AbstractFileSystem:
320320
}
321321

322322
_ADLS_SCHEMES = frozenset({"abfs", "abfss", "wasb", "wasbs"})
323+
_HF_SCHEMES = frozenset({"hf"})
323324

324325

325326
class FsspecInputFile(InputFile):
@@ -351,13 +352,11 @@ def __len__(self) -> int:
351352
def exists(self) -> bool:
352353
"""Check whether the location exists."""
353354
if self._fs_kwargs:
354-
# fsspec's AbstractFileSystem.lexists() doesn't forward **kwargs to exists()/info(),
355-
# so honoring fs_kwargs (e.g. a pinned HuggingFace Hub revision) requires calling info() directly.
356-
try:
357-
self._fs.info(self.location, **self._fs_kwargs)
358-
return True
359-
except FileNotFoundError:
360-
return False
355+
# fsspec's AbstractFileSystem.lexists() doesn't forward **kwargs to exists()/info(), so
356+
# honoring fs_kwargs (e.g. a pinned HuggingFace Hub revision) requires calling exists()
357+
# directly -- it does forward kwargs to info(), with the same broad exception handling
358+
# lexists() would otherwise provide.
359+
return self._fs.exists(self.location, **self._fs_kwargs)
361360
return self._fs.lexists(self.location)
362361

363362
@override
@@ -503,7 +502,7 @@ def _get_fs_kwargs(self, scheme: str) -> Properties:
503502
commit, neither of which is a valid write target, so writes/deletes always target the
504503
repository's default (writable) branch.
505504
"""
506-
if scheme == "hf" and (revision := self.properties.get(HF_REVISION)):
505+
if scheme in _HF_SCHEMES and (revision := self.properties.get(HF_REVISION)):
507506
return {"revision": revision}
508507
return {}
509508

‎tests/io/test_fsspec.py‎

Lines changed: 46 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1134,6 +1134,7 @@ def test_fsspec_hf_revision_forwarded_to_reads() -> None:
11341134
with mock.patch("huggingface_hub.HfFileSystem") as mock_hf_fs:
11351135
mock_fs = mock_hf_fs.return_value
11361136
mock_fs.info.return_value = {"size": 123}
1137+
mock_fs.exists.return_value = True
11371138

11381139
hf_fileio = FsspecFileIO(properties=session_properties)
11391140
input_file = hf_fileio.new_input(location=location)
@@ -1142,7 +1143,7 @@ def test_fsspec_hf_revision_forwarded_to_reads() -> None:
11421143
mock_fs.info.assert_called_with(location, revision="a-pinned-revision")
11431144

11441145
assert input_file.exists() is True
1145-
mock_fs.info.assert_called_with(location, revision="a-pinned-revision")
1146+
mock_fs.exists.assert_called_with(location, revision="a-pinned-revision")
11461147

11471148
input_file.open()
11481149
mock_fs.open.assert_called_with(location, "rb", revision="a-pinned-revision")
@@ -1166,11 +1167,17 @@ def test_fsspec_hf_revision_not_forwarded_to_writes() -> None:
11661167

11671168

11681169
def test_fsspec_hf_no_revision_by_default() -> None:
1169-
with mock.patch("huggingface_hub.HfFileSystem"):
1170+
location = "hf://datasets/user/repo/file.parquet"
1171+
1172+
with mock.patch("huggingface_hub.HfFileSystem") as mock_hf_fs:
1173+
mock_fs = mock_hf_fs.return_value
1174+
mock_fs.info.return_value = {"size": 123}
1175+
11701176
hf_fileio = FsspecFileIO(properties={})
1171-
input_file = hf_fileio.new_input(location="hf://datasets/user/repo/file.parquet")
1177+
input_file = hf_fileio.new_input(location=location)
11721178

1173-
assert input_file._fs_kwargs == {}
1179+
assert len(input_file) == 123
1180+
mock_fs.info.assert_called_with(location)
11741181

11751182

11761183
def test_fsspec_non_hf_scheme_does_not_receive_revision_kwarg() -> None:
@@ -1182,3 +1189,38 @@ def test_fsspec_non_hf_scheme_does_not_receive_revision_kwarg() -> None:
11821189
input_file = fileio.new_input(location="file:///tmp/foo.parquet")
11831190

11841191
assert input_file._fs_kwargs == {}
1192+
1193+
1194+
@pytest.mark.skipif(not os.environ.get("HF_TOKEN"), reason="Requires a real Hugging Face Hub token in HF_TOKEN")
1195+
def test_fsspec_hf_revision_pins_reads_to_a_fixed_commit() -> None:
1196+
"""Without hf.revision, reads always follow the repo's current default branch.
1197+
1198+
This means an Iceberg data-file location pointing into an `hf://` repo isn't reproducible on
1199+
its own: if someone pushes a new commit that changes the same path, every future read of that
1200+
"immutable" data file silently returns the new content instead of what existed when the
1201+
Iceberg snapshot referencing it was written. hf.revision fixes this by letting a table pin
1202+
reads to the exact commit that was current when the file was written.
1203+
"""
1204+
from huggingface_hub import HfApi
1205+
1206+
hf_token = os.environ["HF_TOKEN"]
1207+
api = HfApi(token=hf_token)
1208+
repo_id = f"pyiceberg-hf-revision-test-{uuid.uuid4().hex[:8]}"
1209+
api.create_repo(repo_id=repo_id, repo_type="dataset") # private by default
1210+
location = f"hf://datasets/{repo_id}/file.txt"
1211+
1212+
try:
1213+
first_commit = api.upload_file(path_or_fileobj=b"v1", path_in_repo="file.txt", repo_id=repo_id, repo_type="dataset")
1214+
api.upload_file(path_or_fileobj=b"v2", path_in_repo="file.txt", repo_id=repo_id, repo_type="dataset")
1215+
1216+
# Without hf.revision, reads follow the moving default branch -- the file now reads "v2",
1217+
# even though an Iceberg data file referencing it was "written" back when it was "v1".
1218+
default_branch_fileio = FsspecFileIO(properties={"hf.token": hf_token})
1219+
assert default_branch_fileio.new_input(location).open().read() == b"v2"
1220+
1221+
# Pinning hf.revision to the first commit makes the read reproducible: it keeps returning
1222+
# "v1" regardless of what's since been pushed to the default branch.
1223+
pinned_fileio = FsspecFileIO(properties={"hf.token": hf_token, "hf.revision": first_commit.oid})
1224+
assert pinned_fileio.new_input(location).open().read() == b"v1"
1225+
finally:
1226+
api.delete_repo(repo_id=repo_id, repo_type="dataset")

0 commit comments

Comments
 (0)