From ebb042695a43f203263ec1b5493a909c3a03d157 Mon Sep 17 00:00:00 2001 From: KiriAky 107 Date: Wed, 9 Sep 2026 13:06:28 +0800 Subject: [PATCH] =?UTF-8?q?fix(sync):=20=E6=B5=81=E5=BC=8F=E5=A4=84?= =?UTF-8?q?=E7=90=86=E5=B7=B2=E5=AE=8C=E6=88=90=E5=AF=B9=E8=B1=A1=E4=B8=94?= =?UTF-8?q?=E4=B8=8D=E6=B3=84=E6=BC=8F=E4=BC=A0=E8=BE=93=E8=B5=84=E6=BA=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/acceptance_cases/s05_sync_uploads.py | 17 +++++++++++++++-- server sync/sync_server/storage.py | 15 +++++++++++---- 2 files changed, 26 insertions(+), 6 deletions(-) diff --git a/scripts/acceptance_cases/s05_sync_uploads.py b/scripts/acceptance_cases/s05_sync_uploads.py index a5205bb..c29b8aa 100644 --- a/scripts/acceptance_cases/s05_sync_uploads.py +++ b/scripts/acceptance_cases/s05_sync_uploads.py @@ -115,6 +115,11 @@ def result(case_id: str, status: str, reason: str, facts: dict) -> dict: "used_bytes": facts.get("used_bytes", 0), "reserved_bytes": facts.get("reserved_bytes", -1), }, + { + "scope": "last checkpoint", + "stage": facts.get("active_stage"), + "iteration": facts.get("active_iteration", -1), + }, ], } @@ -236,6 +241,7 @@ def main() -> int: offset_races = durable_offsets = 0 for index in range(ROUNDS): + facts.update({"active_stage": "same-offset", "active_iteration": index}) data = content("offset", index) digest, upload_id, path = begin(base, auth, data) with ThreadPoolExecutor(max_workers=2) as pool: @@ -264,6 +270,7 @@ def main() -> int: disk_ahead = 0 for index in range(ROUNDS): + facts.update({"active_stage": "disk-ahead", "active_iteration": index}) data = content("ahead", index) digest, upload_id, path = begin(base, auth, data) local = stack.staging / upload_id @@ -281,6 +288,7 @@ def main() -> int: disk_behind = 0 for index in range(ROUNDS): + facts.update({"active_stage": "disk-behind", "active_iteration": index}) data = content("behind", index) digest, upload_id, path = begin(base, auth, data) assert call("PUT", path + "?offset=0", headers=auth, body=data[:1])[1] == { @@ -309,6 +317,7 @@ def main() -> int: response_loss_recoveries = 0 for index in range(ROUNDS): + facts.update({"active_stage": "lost-complete", "active_iteration": index}) data = content("lost-complete", index) digest, upload_id, path = begin(base, auth, data) assert call("PUT", path + "?offset=0", headers=auth, body=data)[1] == { @@ -330,6 +339,7 @@ def main() -> int: cleanup_complete = cleanup_removed = 0 cleanup_latencies = [] for index in range(ROUNDS): + facts.update({"active_stage": "expiry-race", "active_iteration": index}) data = content("cleanup", index) digest, upload_id, path = begin(base, auth, data) assert call("PUT", path + "?offset=0", headers=auth, body=data)[1] == { @@ -400,12 +410,14 @@ def main() -> int: assert facts["object_count"] == len(committed) verified = 0 - for digest, data in committed.items(): + for index, (digest, data) in enumerate(committed.items()): + facts.update({"active_stage": "referenced-integrity", "active_iteration": index}) response = call("GET", base + "/objects/" + digest, headers=auth) assert response[0] == 200 and response[1] == data verified += 1 absent = 0 - for digest in cleaned: + for index, digest in enumerate(cleaned): + facts.update({"active_stage": "cleaned-absence", "active_iteration": index}) response = call("GET", base + "/objects/" + digest, headers=auth) assert response[0] == 404 and response[1]["error"]["code"] == "OBJECT_NOT_FOUND" absent += 1 @@ -416,6 +428,7 @@ def main() -> int: assert len(stack.worker_ids()) == 2 assert facts["quota_exact"] and facts["receipt_count_exact"] assert facts["remaining_uploads"] == facts["reserved_bytes"] == facts["staging_files"] == 0 + facts.update({"active_stage": "complete", "active_iteration": ROUNDS}) status = "PASSED" except BaseException as error: reason = "S05_ORACLE_FAILED:" + type(error).__name__ diff --git a/server sync/sync_server/storage.py b/server sync/sync_server/storage.py index 319bf17..c453242 100644 --- a/server sync/sync_server/storage.py +++ b/server sync/sync_server/storage.py @@ -76,11 +76,18 @@ class S3Objects: return stream.read() def put_file(self, key: str, path: Path, content_hash: str): - from boto3.s3.transfer import TransferConfig with path.open("rb") as stream: - self.client.upload_fileobj(stream, self.bucket, key, - ExtraArgs={"Metadata": {"sha256": content_hash}}, - Config=TransferConfig(use_threads=False, max_concurrency=1)) + # Objects are capped at 100 MiB, well below S3's 5 GiB single-PUT + # limit. A direct streaming request has one explicit connection + # lifetime; constructing a transfer manager per completion can + # retain pooled MinIO connections under repeated multi-worker use. + self.client.put_object( + Bucket=self.bucket, + Key=key, + Body=stream, + ContentLength=path.stat().st_size, + Metadata={"sha256": content_hash}, + ) @contextmanager def open(self, key: str):