From a23b61edf11cc2fd0a925c01f3af494218e41b2b Mon Sep 17 00:00:00 2001 From: ASU Date: Wed, 22 Jul 2026 02:59:44 +0300 Subject: [PATCH 1/2] v1.7.0: True seekable MinioBucket --- python/bucketbase/ibucket.py | 3 +- python/bucketbase/minio_bucket.py | 200 +++++++++++++++--- python/bucketbase/versioned_minio_bucket.py | 14 +- python/pyproject.toml | 9 +- python/tests/bucket_tester.py | 54 +++++ python/tests/test_append_only_fs_bucket.py | 5 + python/tests/test_backup_multi_bucket.py | 5 + python/tests/test_fs_bucket.py | 6 + python/tests/test_ibucket.py | 8 + ...test_integrated_cached_immutable_bucket.py | 4 + python/tests/test_memory_bucket.py | 12 ++ python/tests/test_minio_bucket.py | 6 + python/tests/test_versioned_minio_bucket.py | 67 +++++- python/uv.lock | 13 +- 14 files changed, 350 insertions(+), 56 deletions(-) diff --git a/python/bucketbase/ibucket.py b/python/bucketbase/ibucket.py index 178fba3..df386dc 100644 --- a/python/bucketbase/ibucket.py +++ b/python/bucketbase/ibucket.py @@ -249,7 +249,8 @@ def get_object(self, name: PurePosixPath | str) -> bytes: @abstractmethod def get_object_stream(self, name: PurePosixPath | str) -> ObjectStream: """ - Retrieves a stream for reading the object's content. + Retrieves a readable, seekable stream for the object's content. + The yielded stream supports read(), readinto(), seek(), tell(), readable(), and seekable(). :param name: Name of the object to retrieve :return: ObjectStream instance for reading the content diff --git a/python/bucketbase/minio_bucket.py b/python/bucketbase/minio_bucket.py index f33b889..8f45260 100644 --- a/python/bucketbase/minio_bucket.py +++ b/python/bucketbase/minio_bucket.py @@ -2,42 +2,146 @@ import logging import os from pathlib import Path, PurePosixPath +from threading import RLock from typing import BinaryIO, Iterable, Union import certifi import minio import urllib3 +from aicodesign import ai_blackbox from minio import Minio from minio.datatypes import Object from minio.deleteobjects import DeleteError, DeleteObject -from minio.helpers import MAX_PART_SIZE, MIN_PART_SIZE +from minio.helpers import MAX_PART_SIZE, MIN_PART_SIZE, DictType from multiminio import MultiMinio from packaging import version as pkg_version from pyxtension import validate from streamerate import slist from streamerate import stream as sstream -from urllib3 import BaseHTTPResponse +from typing_extensions import override from bucketbase.ibucket import IBucket, ObjectStream, ShallowListing -class MinioObjectStream(ObjectStream): - def __init__(self, response: BaseHTTPResponse, object_name: PurePosixPath) -> None: - # Wrap the BaseHTTPResponse to make it compatible with BinaryIO - super().__init__(response, object_name) # type: ignore[arg-type] - self._response = response - self._size = int(response.headers.get("content-length", -1)) +@ai_blackbox() +class _MinioRangeReader(io.RawIOBase): + @ai_blackbox() + def __init__(self, minio_client: Minio, bucket_name: str, object_name: str, version_id: str | None = None) -> None: + super().__init__() + metadata = minio_client.stat_object(bucket_name, object_name, version_id=version_id) + self._minio_client = minio_client + self._bucket_name = bucket_name + self._object_name = object_name + self._version_id = version_id + if metadata.size is None: + raise IOError(f"Minio returned no size for {object_name}") + self._size = metadata.size + self._etag = metadata.etag + self._position = 0 + self._lock = RLock() + + @override + @ai_blackbox() + def readable(self) -> bool: + self._check_closed() + return True + + @override + @ai_blackbox() + def seekable(self) -> bool: + self._check_closed() + return True + + @override + @ai_blackbox() + def tell(self) -> int: + self._check_closed() + with self._lock: + return self._position + + @override + @ai_blackbox() + def seek(self, offset: int, whence: int = io.SEEK_SET) -> int: + self._check_closed() + with self._lock: + if whence == io.SEEK_SET: + position = offset + elif whence == io.SEEK_CUR: + position = self._position + offset + elif whence == io.SEEK_END: + position = self._size + offset + else: + raise ValueError(f"Invalid whence: {whence}") + if position < 0: + raise ValueError(f"Negative seek position: {position}") + self._position = position + return position + + @override + @ai_blackbox() + def read(self, size: int = -1) -> bytes: + self._check_closed() + with self._lock: + remaining = max(0, self._size - self._position) + length = remaining if size is None or size < 0 else min(size, remaining) + if length == 0: + return b"" + + request_headers: DictType | None = {"If-Match": f'"{self._etag}"'} if self._etag else None + response = self._minio_client.get_object( + self._bucket_name, + self._object_name, + offset=self._position, + length=length, + request_headers=request_headers, + version_id=self._version_id, + ) + try: + data = response.read() + finally: + response.close() + response.release_conn() + + if len(data) != length: + raise IOError(f"Expected {length} bytes from {self._object_name} at offset {self._position}, but received {len(data)}") + self._position += length + return data - def __enter__(self) -> BinaryIO: - return self._response # type: ignore[return-value] + @override + @ai_blackbox() + def readall(self) -> bytes: + return self.read() - def __exit__(self, exc_type: type | None, exc_val: BaseException | None, exc_tb: object | None) -> None: - self._response.close() - self._response.release_conn() + @override + @ai_blackbox() + def readinto(self, buffer: bytearray | memoryview) -> int: + data = self.read(len(buffer)) + buffer[: len(data)] = data + return len(data) + + @ai_blackbox() + def _check_closed(self) -> None: + if self.closed: # pylint: disable=using-constant-test + raise ValueError("I/O operation on closed file") + + +@ai_blackbox() +class MinioObjectStream(ObjectStream): + @ai_blackbox() + def __init__(self, minio_client: Minio, *, bucket_name: str, object_name: PurePosixPath, read_buffer_size: int, version_id: str | None = None) -> None: + raw_reader = _MinioRangeReader(minio_client, bucket_name, str(object_name), version_id=version_id) + stream: BinaryIO = io.BufferedReader(raw_reader, buffer_size=read_buffer_size) if read_buffer_size > 0 else raw_reader + super().__init__(stream, object_name) def build_minio_client( # pylint: disable=too-many-positional-arguments - endpoints: str, access_key: str, secret_key: str, secure: bool = True, region: str | None = "us-east-1", conn_pool_size: int = 128, timeout: int = 5 + endpoints: str, + access_key: str, + secret_key: str, + secure: bool = True, + region: str | None = "us-east-1", + conn_pool_size: int = 128, + timeout: int = 5, ) -> Minio: """ :param endpoints: comma separated list of endpoints @@ -113,43 +217,64 @@ class MinioBucket(IBucket): # Default part size for multipart uploads (16 MiB) # Increased from MinIO library default (5 MiB) for better performance and to allow for objects up to 160GiB DEFAULT_PART_SIZE = 16 * 1024 * 1024 + DEFAULT_READ_BUFFER_SIZE = 128 * 1024 _USES_NAME_ATTRIBUTE: bool = _detect_minio_object_name_attribute() - def __init__(self, bucket_name: str, minio_client: Minio, part_size: int | None = None) -> None: + def __init__( + self, + bucket_name: str, + minio_client: Minio, + part_size: int | None = None, + *, + read_buffer_size: int = DEFAULT_READ_BUFFER_SIZE, + ) -> None: if part_size is None: part_size = self.DEFAULT_PART_SIZE - validate(MIN_PART_SIZE <= part_size <= MAX_PART_SIZE, f"part_size must be between {MIN_PART_SIZE} and {MAX_PART_SIZE}", exc=ValueError) + validate( + MIN_PART_SIZE <= part_size <= MAX_PART_SIZE, + f"part_size must be between {MIN_PART_SIZE} and {MAX_PART_SIZE}", + exc=ValueError, + ) + validate( + isinstance(read_buffer_size, int) and not isinstance(read_buffer_size, bool) and read_buffer_size >= 0, + "read_buffer_size must be a non-negative int", + exc=ValueError, + ) self._minio_client = minio_client self._bucket_name = bucket_name self._part_size = part_size + self._read_buffer_size = read_buffer_size @classmethod def _get_object_name(cls, obj: Object) -> str: - return obj.name if cls._USES_NAME_ATTRIBUTE else obj.object_name + object_name = getattr(obj, "name", None) if cls._USES_NAME_ATTRIBUTE else obj.object_name + if object_name is None: + raise ValueError("Minio object listing item has no object name") + return object_name def get_object(self, name: PurePosixPath | str) -> bytes: with self.get_object_stream(name) as response: - assert isinstance(response, BaseHTTPResponse), f"Expected IOBase, got {type(response)}" - try: - data = bytes() - for buffer in response.stream(amt=1024 * 1024): - data += buffer - return data - finally: - response.release_conn() + return response.read() + + @ai_blackbox() + def _open_object_stream(self, name: str, *, version_id: str | None = None) -> ObjectStream: + return MinioObjectStream( + self._minio_client, + bucket_name=self._bucket_name, + object_name=PurePosixPath(name), + read_buffer_size=self._read_buffer_size, + version_id=version_id, + ) def get_object_stream(self, name: PurePosixPath | str) -> ObjectStream: _name = self._validate_name(name) try: - response: BaseHTTPResponse = self._minio_client.get_object(self._bucket_name, _name) + return self._open_object_stream(_name) except minio.error.S3Error as e: if e.code == "NoSuchKey": raise FileNotFoundError(f"Object {_name} not found in bucket {self._bucket_name} on Minio") from e raise - _name_path = PurePosixPath(_name) if isinstance(_name, str) else _name - return MinioObjectStream(response, _name_path) - def fget_object(self, name: PurePosixPath | str, file_path: Path) -> None: """ Raises: @@ -167,11 +292,22 @@ def put_object(self, name: PurePosixPath | str, content: Union[str, bytes, bytea _content = self._encode_content(content) _name = self._validate_name(name) f = io.BytesIO(_content) - self._minio_client.put_object(bucket_name=self._bucket_name, object_name=_name, data=f, length=len(_content)) + self._minio_client.put_object( + bucket_name=self._bucket_name, + object_name=_name, + data=f, + length=len(_content), + ) def put_object_stream(self, name: PurePosixPath | str, stream: BinaryIO) -> None: _name = self._validate_name(name) - self._minio_client.put_object(bucket_name=self._bucket_name, object_name=_name, data=stream, length=-1, part_size=self._part_size) + self._minio_client.put_object( + bucket_name=self._bucket_name, + object_name=_name, + data=stream, + length=-1, + part_size=self._part_size, + ) def fput_object(self, name: PurePosixPath | str, file_path: Path) -> None: _name = self._validate_name(name) @@ -219,6 +355,8 @@ def remove_objects(self, names: Iterable[PurePosixPath | str]) -> slist[DeleteEr def get_size(self, name: PurePosixPath | str) -> int: try: st = self._minio_client.stat_object(self._bucket_name, str(name)) + if st.size is None: + raise IOError(f"Minio returned no size for {name}") return st.size except minio.error.S3Error as e: if e.code == "NoSuchKey": diff --git a/python/bucketbase/versioned_minio_bucket.py b/python/bucketbase/versioned_minio_bucket.py index 52fa7e2..dc6a1b4 100644 --- a/python/bucketbase/versioned_minio_bucket.py +++ b/python/bucketbase/versioned_minio_bucket.py @@ -7,10 +7,10 @@ from pyxtension import validate from streamerate import slist from streamerate import stream as sstream -from urllib3 import BaseHTTPResponse from bucketbase.ibucket import ObjectStream -from bucketbase.minio_bucket import MinioBucket, MinioObjectStream +from bucketbase.minio_bucket import MinioBucket + @dataclass(frozen=True) class ObjectVersion: @@ -50,25 +50,19 @@ def list_object_versions(self, name: PurePosixPath | str) -> slist[ObjectVersion def get_object_version(self, name: PurePosixPath | str, version_id: str) -> bytes: with self.get_object_version_stream(name, version_id) as response: - assert isinstance(response, BaseHTTPResponse), f"Expected IOBase, got {type(response)}" - data = bytes() - for buffer in response.stream(amt=1024 * 1024): - data += buffer - return data + return response.read() def get_object_version_stream(self, name: PurePosixPath | str, version_id: str) -> ObjectStream: _name = self._validate_name(name) validate(isinstance(version_id, str), f"version_id must be str, but got {type(version_id)}", exc=ValueError) try: - response: BaseHTTPResponse = self._minio_client.get_object(self._bucket_name, _name, version_id=version_id) + return self._open_object_stream(_name, version_id=version_id) except minio.error.S3Error as e: if e.code in ("MethodNotAllowed", "NoSuchKey", "NoSuchVersion"): raise FileNotFoundError(f"Object {_name} version {version_id} not found in bucket {self._bucket_name} on Minio") from e raise - return MinioObjectStream(response, PurePosixPath(_name)) - def remove_object_with_versions(self, name: PurePosixPath | str) -> slist[DeleteError]: versions = self.list_object_versions(name) if versions.size() == 0: diff --git a/python/pyproject.toml b/python/pyproject.toml index fb24028..0230d06 100644 --- a/python/pyproject.toml +++ b/python/pyproject.toml @@ -1,12 +1,13 @@ [project] name = "bucketbase" -version = "1.6.0" # do not edit manually. kept in sync with `tool.commitizen` config via automation +version = "1.7.0" # do not edit manually. kept in sync with `tool.commitizen` config via automation description = "bucketbase" authors = [{ name = "Andrei Suiu", email = "andrei.suiu@gmail.com" }] readme = "README.py.md" license = "MIT" requires-python = ">=3.10,<4.0.0" dependencies = [ + "aicodesign>=0.1.1", "streamerate>=1.2.1,<1.2.7; python_version < '3.11'", "streamerate>=1.2.1; python_version >= '3.11'", "pyxtension>=1.17.1", @@ -40,7 +41,7 @@ dev = [ ] [tool.black] -line-length = 160 +line-length = 180 include = '\.pyi?$' default_language_version = '3.10' @@ -63,7 +64,7 @@ output-format = "colorized" enable = "useless-suppression" [tool.pylint.design] -max-line-length = 160 +max-line-length = 180 max-locals = 25 max-args = 10 @@ -89,7 +90,7 @@ version = "1.2.3" # do not edit manually. kept in sync with `project` config vi tag_format = "v$version" # Same as Black. -line-length = 160 +line-length = 180 [tool.coverage.run] branch = true diff --git a/python/tests/bucket_tester.py b/python/tests/bucket_tester.py index 5059086..bd9a065 100644 --- a/python/tests/bucket_tester.py +++ b/python/tests/bucket_tester.py @@ -14,6 +14,7 @@ import pyarrow as pa import pyarrow.parquet as pq +from aicodesign import ai_blackbox from streamerate import slist from streamerate import stream as sstream @@ -152,6 +153,25 @@ def test_put_and_get_object_stream(self) -> None: result = file.read() self.test_case.assertEqual(result, "Test\ncontent") + @ai_blackbox() + def test_get_object_stream_is_seekable(self) -> None: + path = PurePosixPath(f"dir{self.us}/seekable.bin") + content = b"0123456789abcdefghijklmnopqrstuvwxyz" + self.storage.put_object(path, content) + + with self.storage.get_object_stream(path) as stream: + self.test_case.assertTrue(stream.readable()) + self.test_case.assertTrue(stream.seekable()) + self.test_case.assertEqual(b"0123", stream.read(4)) + self.test_case.assertEqual(10, stream.seek(10)) + self.test_case.assertEqual(b"abcd", stream.read(4)) + self.test_case.assertEqual(12, stream.seek(-2, io.SEEK_CUR)) + self.test_case.assertEqual(b"cd", stream.read(2)) + self.test_case.assertEqual(len(content) - 4, stream.seek(-4, io.SEEK_END)) + destination = bytearray(4) + self.test_case.assertEqual(4, stream.readinto(destination)) # type: ignore[attr-defined] + self.test_case.assertEqual(b"wxyz", destination) + def test_streaming_failure_atomicity(self) -> None: """ Test that failed streaming operations don't leave partial objects in the bucket. @@ -694,6 +714,40 @@ def test_open_write_with_parquet(self) -> None: # pylint: disable=too-many-loca if isinstance(tested_object, AsyncObjectWriter): self.test_case.assertFalse(tested_object._thread.is_alive()) # pylint: disable=W0212 + @ai_blackbox() + def test_get_object_stream_with_parquet_tail(self) -> None: + unique_dir = f"dir{self.us}" + parquet_path = PurePosixPath(f"{unique_dir}/seekable.parquet") + schema = pa.schema([("id", pa.int64()), ("name", pa.string()), ("active", pa.bool_())]) + + parquet_data = io.BytesIO() + with pq.ParquetWriter(parquet_data, schema) as writer: + for row_group in range(3): + start = row_group * 100 + writer.write_batch( + pa.record_batch( + { + "id": list(range(start, start + 100)), + "name": [f"row-{record_id}" for record_id in range(start, start + 100)], + "active": [record_id % 2 == 0 for record_id in range(start, start + 100)], + }, + schema=schema, + ) + ) + self.storage.put_object(parquet_path, parquet_data.getvalue()) + + with self.storage.get_object_stream(parquet_path) as stream: + parquet_file = pq.ParquetFile(stream) + self.test_case.assertTrue(stream.seekable()) + self.test_case.assertEqual(300, parquet_file.metadata.num_rows) + selected = parquet_file.read_row_group(parquet_file.num_row_groups - 1, columns=["id", "active"]) + tail = selected.slice(len(selected) - 75) + + expected_ids = list(range(225, 300)) + self.test_case.assertEqual(["id", "active"], tail.column_names) + self.test_case.assertEqual(expected_ids, tail.column("id").to_pylist()) + self.test_case.assertEqual([record_id % 2 == 0 for record_id in expected_ids], tail.column("active").to_pylist()) + def test_put_object_stream_exception_cleanup(self) -> None: """ Test that when put_object_stream fails (stream raises exception during read), diff --git a/python/tests/test_append_only_fs_bucket.py b/python/tests/test_append_only_fs_bucket.py index c9e4486..956b402 100644 --- a/python/tests/test_append_only_fs_bucket.py +++ b/python/tests/test_append_only_fs_bucket.py @@ -1,3 +1,4 @@ +import io import tempfile import threading import time @@ -121,6 +122,10 @@ def test_get_after_put_object(self): self.assertEqual(retrieved_content, content) obj_stream = bucket_in_test.get_object_stream(object_name) with obj_stream as stream: + self.assertTrue(stream.seekable()) + stream.seek(-7, io.SEEK_END) + self.assertEqual(b"content", stream.read()) + stream.seek(0) retrieved_content_stream_content = stream.read() self.assertEqual(retrieved_content_stream_content, content) diff --git a/python/tests/test_backup_multi_bucket.py b/python/tests/test_backup_multi_bucket.py index d0752d6..ec919b7 100644 --- a/python/tests/test_backup_multi_bucket.py +++ b/python/tests/test_backup_multi_bucket.py @@ -1,5 +1,6 @@ # mypy: disable-error-code="no-untyped-def" import gc +import io import os import tempfile import threading @@ -470,6 +471,10 @@ def test_get_object_stream_functionality(self): self.bucket1.put_object_stream(self.test_name, BytesIO(self.test_content)) with self.multi_bucket.get_object_stream(self.test_name) as stream: + self.assertTrue(stream.seekable()) + stream.seek(-7, io.SEEK_END) + self.assertEqual(self.test_content[-7:], stream.read()) + stream.seek(0) result = stream.read() self.assertEqual(result, self.test_content) diff --git a/python/tests/test_fs_bucket.py b/python/tests/test_fs_bucket.py index 10d3f07..32098de 100644 --- a/python/tests/test_fs_bucket.py +++ b/python/tests/test_fs_bucket.py @@ -31,6 +31,9 @@ def test_put_and_get_object_stream(self) -> None: # check no remaining files in temp dir self.assertEqual(0, len(list((self.storage._root / self.storage.BUCKETBASE_TMP_DIR_NAME).iterdir()))) + def test_get_object_stream_is_seekable(self) -> None: + self.tester.test_get_object_stream_is_seekable() + def test_list_objects(self) -> None: self.tester.test_list_objects() @@ -88,6 +91,9 @@ def test_open_write_feeder_throws(self) -> None: def test_open_write_with_parquet(self) -> None: self.tester.test_open_write_with_parquet() + def test_get_object_stream_with_parquet_tail(self) -> None: + self.tester.test_get_object_stream_with_parquet_tail() + def test_streaming_failure_atomicity(self) -> None: self.tester.test_streaming_failure_atomicity() diff --git a/python/tests/test_ibucket.py b/python/tests/test_ibucket.py index 8b1e8c8..443dd75 100644 --- a/python/tests/test_ibucket.py +++ b/python/tests/test_ibucket.py @@ -170,6 +170,14 @@ def test_open_write_with_parquet(self) -> None: finally: tester.cleanup() + def test_get_object_stream_with_parquet_tail(self) -> None: + bucket = MemoryBucket() + tester = IBucketTester(bucket, self) + try: + tester.test_get_object_stream_with_parquet_tail() + finally: + tester.cleanup() + def test_open_write_timeout(self) -> None: """Test open_write timeout functionality using MemoryBucket.""" bucket = MemoryBucket() diff --git a/python/tests/test_integrated_cached_immutable_bucket.py b/python/tests/test_integrated_cached_immutable_bucket.py index 1164913..22c8ec5 100644 --- a/python/tests/test_integrated_cached_immutable_bucket.py +++ b/python/tests/test_integrated_cached_immutable_bucket.py @@ -89,6 +89,10 @@ def test_get_object_stream_happy_path(self): self.assertFalse(self.cache.exists(path)) self.assertTrue(self.bucket_in_test.exists(path)) with self.bucket_in_test.get_object_stream(path) as stream: + self.assertTrue(stream.seekable()) + stream.seek(-7, io.SEEK_END) + self.assertEqual(b"content", stream.read()) + stream.seek(0) retrieved_content = stream.read() self.assertEqual(retrieved_content, b_content) self.assertTrue(self.bucket_in_test.exists(path)) diff --git a/python/tests/test_memory_bucket.py b/python/tests/test_memory_bucket.py index b7c28f1..5f9fbf2 100644 --- a/python/tests/test_memory_bucket.py +++ b/python/tests/test_memory_bucket.py @@ -19,6 +19,9 @@ def test_put_and_get_object(self): def test_put_and_get_object_stream(self): self.tester.test_put_and_get_object_stream() + def test_get_object_stream_is_seekable(self): + self.tester.test_get_object_stream_is_seekable() + def test_list_objects(self): self.tester.test_list_objects() @@ -49,6 +52,9 @@ def test_open_write_feeder_throws(self): def test_open_write_with_parquet(self): self.tester.test_open_write_with_parquet() + def test_get_object_stream_with_parquet_tail(self): + self.tester.test_get_object_stream_with_parquet_tail() + def test_streaming_failure_atomicity(self): self.tester.test_streaming_failure_atomicity() @@ -130,6 +136,9 @@ def test_put_and_get_object(self): def test_put_and_get_object_stream(self): self.tester.test_put_and_get_object_stream() + def test_get_object_stream_is_seekable(self): + self.tester.test_get_object_stream_is_seekable() + def test_list_objects(self): self.tester.test_list_objects() @@ -160,6 +169,9 @@ def test_open_write_feeder_throws(self): def test_open_write_with_parquet(self): self.tester.test_open_write_with_parquet() + def test_get_object_stream_with_parquet_tail(self): + self.tester.test_get_object_stream_with_parquet_tail() + def test_streaming_failure_atomicity(self): self.tester.test_streaming_failure_atomicity() diff --git a/python/tests/test_minio_bucket.py b/python/tests/test_minio_bucket.py index 91a5ca3..6cc4c8c 100644 --- a/python/tests/test_minio_bucket.py +++ b/python/tests/test_minio_bucket.py @@ -31,6 +31,9 @@ def test_put_and_get_object(self) -> None: def test_put_and_get_object_stream(self) -> None: self.tester.test_put_and_get_object_stream() + def test_get_object_stream_is_seekable(self) -> None: + self.tester.test_get_object_stream_is_seekable() + def test_list_objects(self) -> None: self.tester.test_list_objects() @@ -67,6 +70,9 @@ def test_open_write_feeder_throws(self) -> None: def test_open_write_with_parquet(self) -> None: self.tester.test_open_write_with_parquet() + def test_get_object_stream_with_parquet_tail(self) -> None: + self.tester.test_get_object_stream_with_parquet_tail() + def test_streaming_failure_atomicity(self) -> None: self.tester.test_streaming_failure_atomicity() diff --git a/python/tests/test_versioned_minio_bucket.py b/python/tests/test_versioned_minio_bucket.py index cbe9a57..45a2879 100644 --- a/python/tests/test_versioned_minio_bucket.py +++ b/python/tests/test_versioned_minio_bucket.py @@ -4,6 +4,7 @@ from typing import Iterable, Iterator, Optional from unittest import TestCase +from aicodesign import ai_blackbox from minio import Minio from minio.datatypes import Object from minio.deleteobjects import DeleteError, DeleteObject @@ -53,6 +54,10 @@ def test_full_cycle_object_versions_after_overwrite(self) -> None: self.test_case.assertEqual(b"new content", self.storage.get_object(path)) self.test_case.assertEqual(b"old content", self.storage.get_object_version(path, old_versions[0].version_id)) with self.storage.get_object_version_stream(path, old_versions[0].version_id) as stream: + self.test_case.assertTrue(stream.seekable()) + stream.seek(-7, io.SEEK_END) + self.test_case.assertEqual(b"content", stream.read()) + stream.seek(0) self.test_case.assertEqual(b"old content", stream.read()) errors = self.storage.remove_objects([path]) @@ -95,11 +100,12 @@ def test_invalid_names_raise_for_version_methods(self) -> None: class MockVersionedMinioClient(Minio): def __init__(self) -> None: self.list_objects_response: list[Object] = [] - self.get_object_responses_by_version: dict[str | None, HTTPResponse] = {} + self.content_by_version: dict[str | None, bytes] = {} self.get_object_error: S3Error | None = None self.remove_errors: list[DeleteError] = [] self.list_objects_calls: list[dict[str, object]] = [] self.get_object_calls: list[dict[str, object]] = [] + self.stat_object_calls: list[dict[str, object]] = [] self.remove_objects_calls: list[tuple[str, list[DeleteObject]]] = [] def list_objects(self, **kwargs: object) -> Iterator[Object]: @@ -107,6 +113,25 @@ def list_objects(self, **kwargs: object) -> Iterator[Object]: return iter(self.list_objects_response) @override + @ai_blackbox() + def stat_object( + self, + bucket_name: str, + object_name: str, + ssec: Optional[SseCustomerKey] = None, + version_id: Optional[str] = None, + extra_headers: Optional[DictType] = None, + extra_query_params: Optional[DictType] = None, + ) -> Object: + del ssec, extra_headers, extra_query_params + self.stat_object_calls.append({"bucket_name": bucket_name, "object_name": object_name, "version_id": version_id}) + if self.get_object_error is not None: + raise self.get_object_error + content = self.content_by_version[version_id] + return Object(bucket_name, object_name, etag=f"etag-{version_id}", size=len(content), version_id=version_id) + + @override + @ai_blackbox() def get_object( self, bucket_name: str, @@ -118,10 +143,25 @@ def get_object( version_id: Optional[str] = None, extra_query_params: Optional[DictType] = None, ) -> HTTPResponse: - self.get_object_calls.append({"bucket_name": bucket_name, "object_name": object_name, "version_id": version_id}) + self.get_object_calls.append( + { + "bucket_name": bucket_name, + "object_name": object_name, + "offset": offset, + "length": length, + "request_headers": request_headers, + "version_id": version_id, + } + ) if self.get_object_error is not None: raise self.get_object_error - return self.get_object_responses_by_version[version_id] + content = self.content_by_version[version_id] + return self._make_response(content[offset : offset + length] if length else content[offset:]) + + @staticmethod + @ai_blackbox() + def _make_response(content: bytes) -> HTTPResponse: + return HTTPResponse(body=io.BytesIO(content), headers={"content-length": str(len(content))}, preload_content=False) @override def remove_objects(self, bucket_name: str, delete_object_list: Iterable[DeleteObject], bypass_governance_mode: bool = False) -> Iterator[DeleteError]: @@ -144,10 +184,6 @@ def _make_object(name: str, version_id: str | None, is_latest: str = "false", is is_delete_marker=is_delete_marker, ) - @staticmethod - def _make_response(content: bytes) -> HTTPResponse: - return HTTPResponse(body=io.BytesIO(content), headers={"content-length": str(len(content))}, preload_content=False) - @staticmethod def _make_s3_error(code: str) -> S3Error: return S3Error(HTTPResponse(status=404), code, code, "resource", "request-id", "host-id") @@ -180,12 +216,25 @@ def test_list_object_versions_requires_version_id(self) -> None: self.bucket.list_object_versions("dir/file.txt") def test_get_object_version_reads_specific_version(self) -> None: - self.mock_client.get_object_responses_by_version["v1"] = self._make_response(b"old content") + self.mock_client.content_by_version["v1"] = b"old content" content = self.bucket.get_object_version("dir/file.txt", "v1") self.assertEqual(b"old content", content) - self.assertEqual([{"bucket_name": "test-bucket", "object_name": "dir/file.txt", "version_id": "v1"}], self.mock_client.get_object_calls) + self.assertEqual([{"bucket_name": "test-bucket", "object_name": "dir/file.txt", "version_id": "v1"}], self.mock_client.stat_object_calls) + self.assertEqual( + [ + { + "bucket_name": "test-bucket", + "object_name": "dir/file.txt", + "offset": 0, + "length": len(b"old content"), + "request_headers": {"If-Match": '"etag-v1"'}, + "version_id": "v1", + } + ], + self.mock_client.get_object_calls, + ) def test_get_object_version_requires_string_version_id(self) -> None: with self.assertRaisesRegex(ValueError, "version_id must be str"): diff --git a/python/uv.lock b/python/uv.lock index 34e20e2..fb0588c 100644 --- a/python/uv.lock +++ b/python/uv.lock @@ -8,6 +8,15 @@ resolution-markers = [ "python_full_version < '3.11'", ] +[[package]] +name = "aicodesign" +version = "0.1.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/45/62/97a0fa2256c92903e1de49cec740b6b53dda5f5ed0e6dd047ae40fd064d4/aicodesign-0.1.1.tar.gz", hash = "sha256:8b60ddcff240dad3a5b8c0f55f6f4fea5c92aa4524e50e4bf9bf33983cc526d9", size = 4756, upload-time = "2026-07-18T17:50:14.538Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/6c/d1/6d2b781b2d54fd8876039c6121106e6c057313d0c161f7a6ab84de80f542/aicodesign-0.1.1-py3-none-any.whl", hash = "sha256:caa212232e7354c71be220c824771350f4687a61d11a92cf112fb2eb30b9c59a", size = 3057, upload-time = "2026-07-18T17:50:13.517Z" }, +] + [[package]] name = "annotated-types" version = "0.7.0" @@ -76,9 +85,10 @@ wheels = [ [[package]] name = "bucketbase" -version = "1.6.0" +version = "1.7.0" source = { editable = "." } dependencies = [ + { name = "aicodesign" }, { name = "exceptiongroup", marker = "python_full_version < '3.11'" }, { name = "filelock" }, { name = "pyxtension" }, @@ -110,6 +120,7 @@ dev = [ [package.metadata] requires-dist = [ + { name = "aicodesign", specifier = ">=0.1.1" }, { name = "certifi", marker = "extra == 'minio'", specifier = ">=2024.0.0" }, { name = "exceptiongroup", marker = "python_full_version < '3.11'", specifier = ">=1.0.0" }, { name = "filelock", specifier = ">=3.20.0" }, From 4e4f8bba74ac906ba943f405ffada4ed6ee8b040 Mon Sep 17 00:00:00 2001 From: ASU Date: Wed, 22 Jul 2026 03:49:37 +0300 Subject: [PATCH 2/2] chore: update .gitignore to include .bucketbase.fscache directory --- .gitignore | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/.gitignore b/.gitignore index 643450e..2960135 100644 --- a/.gitignore +++ b/.gitignore @@ -205,4 +205,5 @@ bin/ .vscode/ ### Mac OS ### -.DS_Store \ No newline at end of file +.DS_Store +**/.bucketbase.fscache*/