"""隔离的 PostgreSQL、MinIO 与多工作进程 Sync 验收栈。""" from __future__ import annotations import hashlib import json import os import secrets import socket import subprocess import time from pathlib import Path from urllib.error import HTTPError, URLError from urllib.request import Request, urlopen ROOT = Path(__file__).resolve().parents[2] SERVICE = ROOT / "server sync" MINIO_RELEASE = "RELEASE.2025-09-07T16-13-09Z" MINIO_SHA256 = "af709e6ba68488404e85acdd22a3030d0f5e56a108d4b27d744f18ceb50861b4" def sha256(path: Path) -> str: digest = hashlib.sha256() with path.open("rb") as source: for chunk in iter(lambda: source.read(1024 * 1024), b""): digest.update(chunk) return digest.hexdigest() def free_port() -> int: with socket.socket() as listener: listener.bind(("127.0.0.1", 0)) return listener.getsockname()[1] def call(method: str, url: str, *, headers=None, body=None, timeout=120): request_headers = dict(headers or {}) payload = body if isinstance(body, (dict, list)): payload = json.dumps(body, separators=(",", ":")).encode() request_headers["Content-Type"] = "application/json" request = Request(url, data=payload, headers=request_headers, method=method) try: response = urlopen(request, timeout=timeout) except HTTPError as error: response = error data = response.read() media_type = response.headers.get("Content-Type", "") value = json.loads(data) if data and "json" in media_type else data return response.status, value, {key.lower(): item for key, item in response.headers.items()} def wait_http(url: str, *, expected=200, timeout=30) -> None: until = time.monotonic() + timeout while time.monotonic() < until: try: if call("GET", url)[0] == expected: return except (OSError, URLError): pass time.sleep(0.1) raise RuntimeError("SERVICE_START_TIMEOUT") def stop_tree(process: subprocess.Popen | None) -> None: if process is None or process.poll() is not None: return if os.name == "nt": subprocess.run( ["taskkill", "/PID", str(process.pid), "/T", "/F"], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False, ) else: process.terminate() try: process.wait(timeout=10) except subprocess.TimeoutExpired: process.kill() process.wait(timeout=10) class SyncProductionStack: """拥有一个验收案例的一次性生产依赖堆栈。""" def __init__(self, config: dict, data_root: Path, case_tag: str): self.config = config self.data_root = data_root self.case_tag = case_tag.lower() self.postgres_started = False self.minio = None self.sync = None self.handles = [] self.pg_log = None self.minio_log = None self.sync_log = None self.minio_env = {} initdb = Path(config["artifacts"]["postgres_initdb"]).resolve() self.pg_bin = initdb.parent executable = lambda name: self.pg_bin / (name + ".exe" if os.name == "nt" else name) self.initdb = initdb self.pg_ctl = executable("pg_ctl") self.createdb = executable("createdb") self.postgres = executable("postgres") self.psql = executable("psql") self.minio_server = Path(config["artifacts"]["minio_server"]).resolve() self.server_python = SERVICE / (".venv/Scripts/python.exe" if os.name == "nt" else ".venv/bin/python") required = ( self.initdb, self.pg_ctl, self.createdb, self.postgres, self.psql, self.minio_server, self.server_python, ) if not all(path.is_file() for path in required): raise RuntimeError("PRODUCTION_RUNTIME_MISSING") self.postgres_version = subprocess.check_output( [str(self.postgres), "--version"], text=True, timeout=10 ).strip() self.minio_version = subprocess.check_output( [str(self.minio_server), "--version"], text=True, timeout=10 ).splitlines()[0] if " 17." not in self.postgres_version or MINIO_RELEASE not in self.minio_version: raise RuntimeError("PRODUCTION_RUNTIME_VERSION_MISMATCH") if sha256(self.minio_server) != MINIO_SHA256: raise RuntimeError("MINIO_ARTIFACT_INTEGRITY") self.stack = data_root / f"{self.case_tag}-production-stack" self.postgres_data = self.stack / "postgres" self.minio_data = self.stack / "objects" self.staging = self.stack / "staging" self.pg_port, self.minio_port, self.console_port, self.sync_port = ( free_port() for _ in range(4) ) self.database_url = f"postgresql+psycopg://postgres@127.0.0.1:{self.pg_port}/opennexus" self.origin = f"http://127.0.0.1:{self.sync_port}" self.username = self.case_tag + "-" + secrets.token_hex(8) self.password = secrets.token_urlsafe(32) self.minio_user = self.case_tag + secrets.token_hex(8) self.minio_password = secrets.token_urlsafe(32) self.bucket = "opennexus-" + self.case_tag self.service_env = {} def start(self) -> "SyncProductionStack": self.start_dependencies() self.initialize() self.add_user_with_password(self.username, self.password) self.start_sync() wait_http(self.origin + "/ready", timeout=60) return self def start_dependencies(self) -> "SyncProductionStack": self.stack.mkdir() self.minio_data.mkdir() self.staging.mkdir() self.pg_log = (self.stack / "postgres.log").open("wb") self.minio_log = (self.stack / "minio.log").open("wb") self.sync_log = (self.stack / "sync.log").open("wb") self.handles.extend([self.pg_log, self.minio_log, self.sync_log]) subprocess.run( [ str(self.initdb), "-D", str(self.postgres_data), "-U", "postgres", "-A", "trust", "--no-locale", "-E", "UTF8", ], stdout=self.pg_log, stderr=subprocess.STDOUT, check=True, timeout=120, ) self.start_postgres() subprocess.run( [ str(self.createdb), "-h", "127.0.0.1", "-p", str(self.pg_port), "-U", "postgres", "opennexus", ], stdout=self.pg_log, stderr=subprocess.STDOUT, check=True, timeout=30, ) self.minio_env = os.environ.copy() self.minio_env.update( {"MINIO_ROOT_USER": self.minio_user, "MINIO_ROOT_PASSWORD": self.minio_password} ) self.start_minio() wait_http(f"http://127.0.0.1:{self.minio_port}/minio/health/ready") self.service_env = os.environ.copy() self.service_env.update( { "SYNC_DATABASE_URL": self.database_url, "SYNC_S3_ENDPOINT": f"http://127.0.0.1:{self.minio_port}", "SYNC_S3_BUCKET": self.bucket, "SYNC_STAGING_DIR": str(self.staging), "SYNC_HOST": "127.0.0.1", "SYNC_PORT": str(self.sync_port), "AWS_ACCESS_KEY_ID": self.minio_user, "AWS_SECRET_ACCESS_KEY": self.minio_password, "AWS_DEFAULT_REGION": "us-east-1", "ACCEPTANCE_USERNAME": self.username, "ACCEPTANCE_PASSWORD": self.password, } ) return self def initialize(self) -> None: subprocess.run( [str(self.server_python), "-m", "sync_server", "initialize"], cwd=SERVICE, env=self.service_env, stdout=self.sync_log, stderr=subprocess.STDOUT, check=True, timeout=90, ) def start_sync(self) -> None: if self.sync is not None and self.sync.poll() is None: return self.sync = subprocess.Popen( [str(self.server_python), "-m", "sync_server", "serve", "--workers", "2"], cwd=SERVICE, env=self.service_env, stdout=self.sync_log, stderr=subprocess.STDOUT, ) def stop_sync(self) -> None: stop_tree(self.sync) self.sync = None def login(self, username: str, password: str, label: str) -> dict: session = None for _ in range(5): try: code, candidate, _ = call( "POST", self.origin + "/sync/v1/auth/sessions", body={ "username": username, "password": 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") return session def create_vault_for_auth(self, auth: dict, label: str): 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"] def create_vault(self, label: str): session = self.login(self.username, self.password, label) auth = {"Authorization": "Bearer " + session["access_token"]} return self.create_vault_for_auth(auth, label) def add_user(self, prefix: str) -> tuple[str, str]: username = prefix + "-" + secrets.token_hex(8) password = secrets.token_urlsafe(32) self.add_user_with_password(username, password) return username, password def add_user_with_password(self, username: str, password: str) -> None: environment = dict(self.service_env) environment.update( {"ACCEPTANCE_EXTRA_USERNAME": username, "ACCEPTANCE_EXTRA_PASSWORD": password} ) subprocess.run( [ str(self.server_python), "-c", "import os; from sync_server.database import Database; " "Database(os.environ['SYNC_DATABASE_URL']).add_user(" "os.environ['ACCEPTANCE_EXTRA_USERNAME'],os.environ['ACCEPTANCE_EXTRA_PASSWORD'])", ], cwd=SERVICE, env=environment, stdout=self.sync_log, stderr=subprocess.STDOUT, check=True, timeout=30, ) def start_postgres(self) -> None: subprocess.run( [ str(self.pg_ctl), "-D", str(self.postgres_data), "-l", str(self.stack / "postgres-server.log"), "-o", f"-p {self.pg_port} -h 127.0.0.1", "-w", "start", ], stdout=self.pg_log, stderr=subprocess.STDOUT, check=True, timeout=60, ) self.postgres_started = True def stop_postgres(self) -> None: if not self.postgres_started: return subprocess.run( [str(self.pg_ctl), "-D", str(self.postgres_data), "-m", "fast", "-w", "stop"], stdout=self.pg_log, stderr=subprocess.STDOUT, check=True, timeout=30, ) self.postgres_started = False def start_minio(self) -> None: if self.minio is not None and self.minio.poll() is None: return self.minio = subprocess.Popen( [ str(self.minio_server), "server", str(self.minio_data), "--address", f"127.0.0.1:{self.minio_port}", "--console-address", f"127.0.0.1:{self.console_port}", ], stdout=self.minio_log, stderr=subprocess.STDOUT, env=self.minio_env, ) def stop_minio(self) -> None: stop_tree(self.minio) self.minio = None def sql_scalar(self, statement: str) -> int: output = subprocess.check_output( [ str(self.psql), "-X", "-A", "-t", "-v", "ON_ERROR_STOP=1", "-h", "127.0.0.1", "-p", str(self.pg_port), "-U", "postgres", "-d", "opennexus", "-c", statement, ], text=True, timeout=30, ).strip() return int(output) def worker_ids(self, attempts=200) -> set[str]: from concurrent.futures import ThreadPoolExecutor def probe(_): try: return call( "GET", self.origin + "/health", headers={"Connection": "close"}, timeout=10, )[2].get("x-opennexus-worker") except (OSError, TimeoutError, URLError): return None # 一个小池足以覆盖两个工作线程,而不会在故障矩阵启动之前耗尽 Windows 临时端口/积压预算。 with ThreadPoolExecutor(max_workers=8) as pool: return {value for value in pool.map(probe, range(attempts)) if value} def assert_running(self) -> None: if self.sync is None or self.minio is None: raise RuntimeError("PRODUCTION_STACK_NOT_STARTED") if self.sync.poll() is not None or self.minio.poll() is not None: raise RuntimeError("PRODUCTION_STACK_EXITED") def stop(self) -> None: self.stop_sync() self.stop_minio() if self.postgres_started: try: self.stop_postgres() except (OSError, subprocess.SubprocessError): pass for handle in self.handles: handle.close()