From 4744c5888e543928070392a29b7337a6b84d7074 Mon Sep 17 00:00:00 2001 From: hannes Date: Tue, 28 Jul 2026 22:12:39 +0200 Subject: [PATCH] fix: retry UploadPart on transient backend errors RGW can accept a part, commit to a 200, stream it, and only then put the failure in the body ("InternalError: The server did not respond in time.", status code: 200). Three retry layers all miss that shape because each decides from the HTTP status line, which already read 200: - rclone's LowLevelRetries gates on status 429/500/503 (backend/s3/s3.go) - rclone's string fallback matches only transport phrases ("use of closed network connection", "unexpected EOF reading trailer"), not InternalError - botocore's max_attempts in S3Client sees the same 200 #133 added this retry to UploadPartCopy and #138 to CompleteMultipartUpload, but plain UploadPart was left bare -- and it is the only one of the three a Scylla backup uses (267 PUT, 0 UploadPartCopy measured against the backup bucket). At 50MB rclone chunks a 6.2GB SSTable is ~762 internal PUTs, so one unretried transient failure loses the whole multi-GB file. 2026-07-28 prod: 10 of 13 racks failed, ~2.4TiB never uploaded, every failure on main.companies (47 of 48 large files). Reuses the existing base.SOURCE_READ_ATTEMPTS machinery from #133, so no new tunables. Safe at both call sites: the full ciphertext is still in memory when upload_part is called (the del happens after), so a retry re-PUTs identical bytes to the same internal part number. --- s3proxy/handlers/multipart/upload_part.py | 57 +++++++++- tests/unit/test_upload_part_retry.py | 122 ++++++++++++++++++++++ 2 files changed, 176 insertions(+), 3 deletions(-) create mode 100644 tests/unit/test_upload_part_retry.py diff --git a/s3proxy/handlers/multipart/upload_part.py b/s3proxy/handlers/multipart/upload_part.py index cfc3048..02e8561 100644 --- a/s3proxy/handlers/multipart/upload_part.py +++ b/s3proxy/handlers/multipart/upload_part.py @@ -25,7 +25,8 @@ StateMissingError, ) from ...streaming import decode_aws_chunked_stream -from ..base import BaseHandler +from .. import base +from ..base import BaseHandler, is_retryable_source_error logger: BoundLogger = structlog.get_logger(__name__) @@ -236,6 +237,46 @@ async def handle_upload_part(self, request: Request, creds: S3Credentials) -> Re except Exception as e: return self._handle_generic_error(e, bucket, key, part_num, upload_id) + async def _upload_part_with_retry( + self, + client: S3Client, + bucket: str, + key: str, + upload_id: str, + internal_part_num: int, + ciphertext: bytes | bytearray, + *, + client_part_num: int, + ) -> dict: + """UploadPart, retrying transient backend failures. + + RGW can accept the part, commit to a 200, and only then fail with an + body ("InternalError: The server did not respond in time"). + Neither botocore's retries nor rclone's low-level retries fire for that: + both decide from the HTTP status line, which already read 200. Re-PUTting + the same part number with the same ciphertext is idempotent, so retry here + where the parsed error is actually visible. + """ + for attempt in range(1, base.SOURCE_READ_ATTEMPTS + 1): + try: + return await client.upload_part( + bucket, key, upload_id, internal_part_num, ciphertext + ) + except Exception as exc: + if not is_retryable_source_error(exc) or attempt == base.SOURCE_READ_ATTEMPTS: + raise + logger.warning( + "UPLOAD_PART_RETRY", + bucket=bucket, + key=key, + client_part=client_part_num, + internal_part=internal_part_num, + attempt=attempt, + error=str(exc), + ) + await asyncio.sleep(base.SOURCE_READ_BACKOFF_SEC * (2 ** (attempt - 1))) + raise AssertionError("unreachable") + async def _get_or_recover_state( self, client: S3Client, bucket: str, key: str, upload_id: str, part_num: int ) -> MultipartUploadState: @@ -462,7 +503,9 @@ async def _stream_and_upload_framed( ) upload_start = time.monotonic() - resp = await client.upload_part(bucket, key, upload_id, ipn, ciphertext) + resp = await self._upload_part_with_retry( + client, bucket, key, upload_id, ipn, ciphertext, client_part_num=part_num + ) del ciphertext etag = resp["ETag"].strip('"') logger.info( @@ -603,7 +646,15 @@ async def _upload_internal_part_with_semaphore( del data # Free memory # Upload - resp = await client.upload_part(bucket, key, upload_id, internal_part_num, ciphertext) + resp = await self._upload_part_with_retry( + client, + bucket, + key, + upload_id, + internal_part_num, + ciphertext, + client_part_num=client_part_num, + ) etag = resp["ETag"].strip('"') del ciphertext # Free memory diff --git a/tests/unit/test_upload_part_retry.py b/tests/unit/test_upload_part_retry.py new file mode 100644 index 0000000..dbe4bb8 --- /dev/null +++ b/tests/unit/test_upload_part_retry.py @@ -0,0 +1,122 @@ +"""UploadPart retry on transient backend errors. + +Prod incident 2026-07-28: the Scylla backup of `main.companies` failed on 10 of +13 racks, leaving ~2.4TiB un-uploaded. Every failure was a ~6.2GiB SSTable and +every error was ``InternalError: The server did not respond in time.`` with +``status code: 200`` -- RGW commits to a 200, streams it, then puts the failure +in the body. + +Three retry layers all miss that shape because each decides from the HTTP status +line, which already read 200: + * rclone's ``LowLevelRetries`` gates on status 429/500/503 (backend/s3/s3.go) + * rclone's string fallback only matches transport phrases ("use of closed + network connection", "unexpected EOF reading trailer"), not ``InternalError`` + * botocore's ``max_attempts`` in S3Client sees the same 200 + +#133 added this retry to UploadPartCopy and #138 to CompleteMultipartUpload, but +UploadPart -- the operation a Scylla backup actually uses (all PUT, no COPY) -- +was left bare. At 50MB rclone chunks a 6.2GiB file is ~762 internal PUTs, so one +unretried transient failure loses the whole multi-GB file. +""" + +from __future__ import annotations + +import pytest +from botocore.exceptions import ClientError + +from s3proxy.handlers import base +from s3proxy.handlers.multipart.upload_part import UploadPartMixin + + +@pytest.fixture(autouse=True) +def _no_backoff(monkeypatch): + monkeypatch.setattr(base, "SOURCE_READ_BACKOFF_SEC", 0.0) + + +def _client_error(code: str, message: str | None = None) -> ClientError: + return ClientError({"Error": {"Code": code, "Message": message or code}}, "UploadPart") + + +class _FlakyUploadClient: + """Backend whose first ``fail_times`` upload_part calls raise ``error_code``.""" + + def __init__(self, fail_times: int = 0, error_code: str = "InternalError"): + self._fail_times = fail_times + self._error_code = error_code + self.attempts = 0 + self.seen = [] + + async def upload_part(self, bucket, key, upload_id, part_number, body): + self.attempts += 1 + self.seen.append((part_number, bytes(body))) + if self._fail_times > 0: + self._fail_times -= 1 + if self._error_code == "InternalError": + raise _client_error("InternalError", "The server did not respond in time.") + raise _client_error(self._error_code) + return {"ETag": f'"etag-{part_number}"'} + + +async def _upload(client, **kw): + return await UploadPartMixin._upload_part_with_retry( + UploadPartMixin.__new__(UploadPartMixin), + client, + "bucket", + "key", + "upload-id", + kw.pop("internal_part_num", 7), + kw.pop("ciphertext", b"ciphertext-bytes"), + client_part_num=kw.pop("client_part_num", 1), + ) + + +@pytest.mark.asyncio +async def test_retries_late_200_internal_error(): + """The exact prod shape: InternalError in a 200 body, then success.""" + client = _FlakyUploadClient(fail_times=1) + resp = await _upload(client) + assert resp["ETag"] == '"etag-7"' + assert client.attempts == 2 + + +@pytest.mark.asyncio +async def test_retry_reuses_same_part_number_and_bytes(): + """Re-PUT must be idempotent - same part number, same ciphertext.""" + client = _FlakyUploadClient(fail_times=2) + await _upload(client, internal_part_num=3, ciphertext=b"abc") + assert client.attempts == 3 + assert client.seen == [(3, b"abc")] * 3 + + +@pytest.mark.asyncio +async def test_gives_up_after_attempt_cap(): + """Exhausting attempts re-raises rather than looping forever.""" + client = _FlakyUploadClient(fail_times=base.SOURCE_READ_ATTEMPTS) + with pytest.raises(ClientError): + await _upload(client) + assert client.attempts == base.SOURCE_READ_ATTEMPTS + + +@pytest.mark.asyncio +async def test_non_retryable_error_is_not_retried(): + """A real client error must fail fast, not burn attempts.""" + client = _FlakyUploadClient(fail_times=1, error_code="AccessDenied") + with pytest.raises(ClientError): + await _upload(client) + assert client.attempts == 1 + + +@pytest.mark.asyncio +async def test_success_path_makes_one_call(): + client = _FlakyUploadClient() + resp = await _upload(client) + assert resp["ETag"] == '"etag-7"' + assert client.attempts == 1 + + +@pytest.mark.parametrize("code", ["InternalError", "SlowDown", "ServiceUnavailable"]) +@pytest.mark.asyncio +async def test_retryable_codes(code): + client = _FlakyUploadClient(fail_times=1, error_code=code) + await _upload(client) + assert client.attempts == 2