From 2858babfe3f3e70196e5ea0ecaf6b19ff75cb563 Mon Sep 17 00:00:00 2001 From: KiriAky 107 Date: Wed, 9 Sep 2026 12:56:33 +0800 Subject: [PATCH] =?UTF-8?q?test(sync):=20=E7=A8=B3=E5=AE=9A=E7=94=9F?= =?UTF-8?q?=E4=BA=A7=E6=A0=88=E6=8E=A2=E6=B5=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/acceptance_cases/s05_sync_uploads.py | 4 +- .../acceptance_cases/sync_production_stack.py | 67 ++++++++++++++----- 2 files changed, 51 insertions(+), 20 deletions(-) diff --git a/scripts/acceptance_cases/s05_sync_uploads.py b/scripts/acceptance_cases/s05_sync_uploads.py index 3f065a1..a5205bb 100644 --- a/scripts/acceptance_cases/s05_sync_uploads.py +++ b/scripts/acceptance_cases/s05_sync_uploads.py @@ -225,11 +225,11 @@ def main() -> int: "dependencies": True, "postgres_version": stack.postgres_version, "minio_version": stack.minio_version, - "worker_count": len(stack.worker_ids()), } ) - assert facts["worker_count"] == 2 auth, base, vault_id = stack.create_vault("S-05 isolated") + facts["worker_count"] = len(stack.worker_ids()) + assert facts["worker_count"] == 2 committed: dict[str, bytes] = {} cleaned: dict[str, bytes] = {} expected_used = 0 diff --git a/scripts/acceptance_cases/sync_production_stack.py b/scripts/acceptance_cases/sync_production_stack.py index f9400bb..d91c9b5 100644 --- a/scripts/acceptance_cases/sync_production_stack.py +++ b/scripts/acceptance_cases/sync_production_stack.py @@ -267,22 +267,45 @@ class SyncProductionStack: return self def create_vault(self, label: str): - code, session, _ = call( - "POST", - self.origin + "/sync/v1/auth/sessions", - body={ - "username": self.username, - "password": self.password, - "device_name": label, - }, - ) - if code != 200: + session = None + for _ in range(5): + try: + code, candidate, _ = call( + "POST", + self.origin + "/sync/v1/auth/sessions", + body={ + "username": self.username, + "password": self.password, + "device_name": label, + }, + timeout=10, + ) + if code == 200: + session = candidate + break + except (OSError, TimeoutError, URLError): + pass + time.sleep(0.2) + if session is None: raise RuntimeError("ACCEPTANCE_LOGIN_FAILED") auth = {"Authorization": "Bearer " + session["access_token"]} - code, vault, _ = call( - "POST", self.origin + "/sync/v1/vaults", headers=auth, body={"name": label} - ) - if code != 200: + vault = None + for _ in range(5): + try: + code, candidate, _ = call( + "POST", + self.origin + "/sync/v1/vaults", + headers=auth, + body={"name": label}, + timeout=10, + ) + if code == 200: + vault = candidate + break + except (OSError, TimeoutError, URLError): + pass + time.sleep(0.2) + if vault is None: raise RuntimeError("ACCEPTANCE_VAULT_FAILED") return auth, self.origin + "/sync/v1/vaults/" + vault["vault_id"], vault["vault_id"] @@ -315,11 +338,19 @@ class SyncProductionStack: from concurrent.futures import ThreadPoolExecutor def probe(_): - return call( - "GET", self.origin + "/health", headers={"Connection": "close"} - )[2].get("x-opennexus-worker") + try: + return call( + "GET", + self.origin + "/health", + headers={"Connection": "close"}, + timeout=10, + )[2].get("x-opennexus-worker") + except (OSError, TimeoutError, URLError): + return None - with ThreadPoolExecutor(max_workers=40) as pool: + # A small pool is enough to reach both workers without exhausting the + # Windows ephemeral-port/backlog budget before the fault matrix starts. + with ThreadPoolExecutor(max_workers=8) as pool: return {value for value in pool.map(probe, range(attempts)) if value} def assert_running(self) -> None: