Skip to content

Commit 130217c

Browse files
Merge branch 'develop' into fix/dynamodb-number-magnitude
2 parents b0255b1 + 3311cc1 commit 130217c

2 files changed

Lines changed: 101 additions & 10 deletions

File tree

‎aws_lambda_powertools/utilities/streaming/_s3_seekable_io.py‎

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,10 @@
55
from typing import IO, TYPE_CHECKING, Any, TypeVar
66

77
import boto3
8+
from botocore.exceptions import ClientError
89

910
from aws_lambda_powertools.shared import user_agent
11+
from aws_lambda_powertools.utilities.streaming.compat import PowertoolsStreamingBody
1012
from aws_lambda_powertools.utilities.streaming.constants import MESSAGE_STREAM_NOT_WRITABLE
1113

1214
if TYPE_CHECKING:
@@ -15,8 +17,6 @@
1517

1618
from mypy_boto3_s3.client import S3Client
1719

18-
from aws_lambda_powertools.utilities.streaming.compat import PowertoolsStreamingBody
19-
2020
_CData = TypeVar("_CData")
2121

2222
logger = logging.getLogger(__name__)
@@ -98,10 +98,21 @@ def raw_stream(self) -> PowertoolsStreamingBody:
9898
"""
9999
Returns the boto3 StreamingBody, starting the stream from the sought position.
100100
"""
101+
if self._closed:
102+
raise ValueError("I/O operation on closed file.")
103+
101104
if self._raw_stream is None:
102105
range_header = f"bytes={self._position}-"
103106
logger.debug(f"Starting new stream at {range_header}")
104-
self._raw_stream = self.s3_client.get_object(Range=range_header, **self._sdk_options).get("Body")
107+
try:
108+
self._raw_stream = self.s3_client.get_object(Range=range_header, **self._sdk_options).get("Body")
109+
except ClientError as exc:
110+
# S3 rejects a range that starts at or past the end of the object, which includes any range
111+
# on an empty object. A file returns no data at that position instead of raising.
112+
if exc.response.get("Error", {}).get("Code") != "InvalidRange":
113+
raise
114+
logger.debug(f"Position {self._position} is at or past the end of the object")
115+
self._raw_stream = PowertoolsStreamingBody(raw_stream=io.BytesIO(b""), content_length=0)
105116
self._closed = False
106117

107118
return self._raw_stream
@@ -183,7 +194,9 @@ def __exit__(self, *kwargs):
183194
self.close()
184195

185196
def close(self) -> None:
186-
self.raw_stream.close()
197+
# Only close a stream that is already open, rather than opening a new one just to close it
198+
if self._raw_stream is not None:
199+
self._raw_stream.close()
187200
self._closed = True
188201

189202
def fileno(self) -> int:

‎tests/functional/streaming/_boto3/test_s3_seekable_io.py‎

Lines changed: 84 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,9 @@
55
import boto3
66
import pytest
77
from botocore import stub
8+
from botocore.exceptions import ClientError
89

10+
from aws_lambda_powertools.utilities.streaming import S3Object
911
from aws_lambda_powertools.utilities.streaming._s3_seekable_io import _S3SeekableIO
1012
from aws_lambda_powertools.utilities.streaming.compat import PowertoolsStreamingBody
1113

@@ -134,17 +136,93 @@ def test_readlines(s3_seekable_obj, s3_client_stub):
134136
assert s3_seekable_obj.tell() == len(payload)
135137

136138

137-
def test_closed(s3_seekable_obj, s3_client_stub):
138-
payload = b"test"
139-
streaming_body = PowertoolsStreamingBody(raw_stream=io.BytesIO(payload), content_length=len(payload))
139+
def test_read_at_end_of_object_returns_empty_bytes(s3_seekable_obj, s3_client_stub):
140+
s3_client_stub.add_response("head_object", {"ContentLength": 4})
141+
# S3 rejects a range that starts at the end of the object
142+
s3_client_stub.add_client_error(
143+
"get_object",
144+
service_error_code="InvalidRange",
145+
http_status_code=416,
146+
expected_params={"Bucket": s3_seekable_obj.bucket, "Key": s3_seekable_obj.key, "Range": "bytes=4-"},
147+
)
140148

141-
s3_client_stub.add_response(
149+
s3_seekable_obj.seek(0, io.SEEK_END)
150+
151+
assert s3_seekable_obj.read() == b""
152+
assert s3_seekable_obj.tell() == 4
153+
154+
155+
def test_read_empty_object_returns_empty_bytes(s3_seekable_obj, s3_client_stub):
156+
# S3 rejects any range on an empty object
157+
s3_client_stub.add_client_error(
142158
"get_object",
143-
{"Body": streaming_body},
144-
{"Bucket": s3_seekable_obj.bucket, "Key": s3_seekable_obj.key, "Range": "bytes=0-"},
159+
service_error_code="InvalidRange",
160+
http_status_code=416,
161+
expected_params={"Bucket": s3_seekable_obj.bucket, "Key": s3_seekable_obj.key, "Range": "bytes=0-"},
145162
)
146163

164+
assert s3_seekable_obj.read() == b""
165+
assert list(s3_seekable_obj) == []
166+
assert s3_seekable_obj.tell() == 0
167+
168+
169+
def test_raw_stream_raises_other_client_errors(s3_seekable_obj, s3_client_stub):
170+
s3_client_stub.add_client_error("get_object", service_error_code="NoSuchKey", http_status_code=404)
171+
172+
with pytest.raises(ClientError, match="NoSuchKey"):
173+
s3_seekable_obj.read()
174+
175+
176+
def test_closed(s3_seekable_obj, s3_client_stub):
147177
s3_seekable_obj.close()
178+
179+
assert s3_seekable_obj.closed is True
180+
# Closing an object that was never read must not open a stream just to close it
181+
s3_client_stub.assert_no_pending_responses()
182+
183+
184+
@pytest.mark.parametrize("stream_class", [_S3SeekableIO, S3Object])
185+
@pytest.mark.parametrize("read_method", ["read", "readline", "readlines", "__next__"])
186+
@pytest.mark.parametrize("initial_state", ["unread", "partially_read", "seeked", "empty"])
187+
def test_reads_after_close_do_not_reopen_stream(s3_client, s3_client_stub, stream_class, read_method, initial_state):
188+
stream = stream_class(bucket="bucket", key="key", boto3_client=s3_client)
189+
expected_params = {"Bucket": "bucket", "Key": "key", "Range": "bytes=0-"}
190+
191+
if initial_state == "empty":
192+
s3_client_stub.add_client_error(
193+
"get_object",
194+
service_error_code="InvalidRange",
195+
http_status_code=416,
196+
expected_params=expected_params,
197+
)
198+
assert stream.read() == b""
199+
elif initial_state != "unread":
200+
payload = b"hello\nworld"
201+
body = PowertoolsStreamingBody(raw_stream=io.BytesIO(payload), content_length=len(payload))
202+
s3_client_stub.add_response("get_object", {"Body": body}, expected_params)
203+
assert stream.read(1) == b"h"
204+
if initial_state == "seeked":
205+
stream.seek(3)
206+
207+
position = stream.tell()
208+
stream.close()
209+
stream.close()
210+
211+
with pytest.raises(ValueError, match="I/O operation on closed file"):
212+
getattr(stream, read_method)()
213+
214+
assert stream.closed is True
215+
assert stream.tell() == position
216+
# The stub has no queued responses, so any attempt to reopen the stream would fail the test.
217+
s3_client_stub.assert_no_pending_responses()
218+
219+
220+
def test_context_manager_at_end_of_object(s3_seekable_obj, s3_client_stub):
221+
s3_client_stub.add_response("head_object", {"ContentLength": 4})
222+
223+
with s3_seekable_obj as f:
224+
f.seek(0, io.SEEK_END)
225+
148226
assert s3_seekable_obj.closed is True
149227

150228

0 commit comments

Comments
 (0)