test(sync): 验证四个并发大型 HTTP 上传
This commit is contained in:
@@ -652,3 +652,12 @@ Core 的独立数据目录目前不等于已授权 Vault。Python 旧笔记写
|
||||
- upload_chunk 原先在 async 路由中直接执行同步数据库锁等待、文件写入和 fsync,会阻塞该 worker 事件循环。现将整段事务/文件持久化交给 Starlette run_in_threadpool,连接在同一工作线程创建和提交,仍保持先刷盘再确认 offset 的协议顺序。接收块先检查剩余额度再扩展 bytearray,避免把超限块复制进应用缓冲区;ASGI 自身收到的块不计入该应用缓冲上限证明。
|
||||
- 新增故障测试以 Event 持续阻塞 fsync,在同一 TestClient ASGI 事件循环中要求 /health 在 2 秒内响应且上传仍在等待;释放后检查偏移与 complete。另测 1 MiB 加 1 字节拒绝、fsync 异常事务回滚、未确认尾部协调截断、恰好 1 MiB 重试成功。
|
||||
- uv run --project "server sync" pytest "server sync/tests" -q 全套 25 通过,耗时 7.80 秒,JUnit:.build/sync-server-upload-threadpool.xml。依赖报告 2 项 TestClient/AnyIO 弃用警告,无测试失败。环境为 SQLite/DiskObjects 受控测试,不是 PostgreSQL/MinIO 两 worker 或 S-09 四并发 100 MiB 服务总 RSS/30 分钟负载证明;线程池饱和与真实部署验收仍需继续。完整生产化目标未完成。
|
||||
|
||||
|
||||
## 增量:四客户端大附件传输验收入口
|
||||
|
||||
- 新增 server sync/tools/upload_benchmark.py,固定种子 20260908,两测试账号各两个新 Vault,四任务屏障同时开始默认 100 MiB 上传。内容逐个 1 MiB 生成,下载也逐块验证,不在驱动器构造完整附件。逐 offset、完整回执与 revision 重放、SHA256/长度、跨账号对象拒绝均为必过条件;准备失败会取消其他等待任务。
|
||||
- 命令入口从本地凭据文件读取两个测试账号,不把凭据/令牌/异常原文写入输出;拒绝带认证信息/query/fragment 的 origin,不跟随重定向,HTTP 要求显式测试开关。整体运行限时 30 分钟;保留四个新测试 Vault,不删除既有数据。README 已列命令与隔离数据要求。
|
||||
- 实际运行 Uvicorn TCP 回环测试:四个 1 MiB + 17 字节验证非整块末尾;四个 100 MiB 验证实际大对象。两项通过,耗时 16.28 秒;大对象上传至下载/隔离验证全部完成约 9.31 秒,400 MiB 内容全部摘要一致。日志 .build/sync-upload-benchmark.log,JUnit .build/sync-upload-benchmark.xml。
|
||||
- 此结果是同一 Python 进程内 SQLite/DiskObjects 单服务实例与 HTTP 客户端的受控检查,不是 PostgreSQL/MinIO 两 worker 基准。service_rss_bytes 明确 null,acceptance 明确 NOT_ASSESSED,不能将传输通过当作 RSS≤2 GiB 或完整 S-09 通过。本机 Get-Command docker 未找到可执行文件;生产拓扑/服务内存采集/30 分钟负载仍需继续。
|
||||
- 最终 Sync 服务端全套 27 通过、2 项依赖弃用警告,22.39 秒;日志 .build/sync-server-four-upload-full.log,JUnit .build/sync-server-four-upload-full.xml。用户 Vault 修改保持原状,完整生产化目标未完成。
|
||||
|
||||
@@ -25,3 +25,16 @@ Dockerfile/Compose 是待实测部署配置,不能视为已验证安装程序
|
||||
2026-09-08 已在独立 Docker 项目完成真实 PostgreSQL/MinIO 双 worker 测试部署,修正基础镜像中的 `sync` 系统用户名冲突。测试专用 HTTP 地址、故障检查、完整验收缺口与运维入口见[验收报告](../docs/development/OpenNexus验收报告-2026-09-08.md)。仓库通用 Compose 的生产 TLS 与备份恢复仍未通过发布验收。
|
||||
|
||||
升级前同时备份 PostgreSQL 和 Bucket,停止提交以取得一致切点。Schema v1 拒绝未知数据库版本,不自动降级。当前历史永久保留,容量管理不能手动删除被历史引用的对象。
|
||||
|
||||
|
||||
## 四并发大附件传输探针
|
||||
|
||||
仅对可保留测试数据的隔离实例运行。准备两个已有测试账号的受限权限 JSON 文件,内容为两个含 username/password 字段的对象数组;不要将该文件提交到仓库。工具会创建四个测试 Vault,保留数据供进一步核对,不修改已有 Vault。
|
||||
|
||||
```powershell
|
||||
uv run python tools/upload_benchmark.py --url http://localhost:8080 --allow-test-http --credentials C:/private/sync-test-accounts.json --output ../.build/upload-report.json
|
||||
```
|
||||
|
||||
默认两账号各两个 Vault,以共同启动屏障进行四个 100 MiB 上传,逐块确认 offset,重复 complete 和 revision 检查持久回执,再流式下载核对大小/SHA256并检查跨账号拒绝。无自动重试掩盖失败,报告不含凭据、令牌或响应正文。HTTP 必须显式指定测试选项;HTTPS 使用默认验证且不跟随重定向。单次运行上限 30 分钟。
|
||||
|
||||
报告 result 表示本次传输检查结果,acceptance 始终 NOT_ASSESSED:工具尚未接入服务进程/容器 RSS 采集器,也未执行完整 S-09 的 30 分钟提交负载、10000 笔记和规定网络条件。真实 TCP 回环回归见 tests/test_upload_benchmark.py,其中 SQLite/DiskObjects 与客户端位于同一 Python 进程,不能作为生产性能指标。
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
"""Real localhost HTTP harness; SQLite/DiskObjects, not production topology."""
|
||||
import asyncio
|
||||
import json
|
||||
import socket
|
||||
|
||||
import pytest
|
||||
import uvicorn
|
||||
|
||||
from sync_server.app import create_app
|
||||
from sync_server.database import Database
|
||||
from sync_server.storage import DiskObjects
|
||||
from tools.upload_benchmark import probe
|
||||
|
||||
|
||||
@pytest.mark.parametrize("size", [1048576 + 17, 100 * 1048576])
|
||||
def test_four_concurrent_uploads_over_real_http(tmp_path, size):
|
||||
database = Database("sqlite:///" + str(tmp_path / "sync.db"))
|
||||
app = create_app(database, DiskObjects(tmp_path / "objects"), tmp_path / "staging")
|
||||
credentials = [{"username": username, "password": "controlled-fixture-password"}
|
||||
for username in ("benchmark-one", "benchmark-two")]
|
||||
for account in credentials:
|
||||
database.add_user(**account)
|
||||
sock = socket.socket()
|
||||
sock.bind(("127.0.0.1", 0))
|
||||
sock.listen(128)
|
||||
server = uvicorn.Server(uvicorn.Config(app, log_config=None, access_log=False,
|
||||
timeout_graceful_shutdown=2))
|
||||
|
||||
async def run():
|
||||
task = asyncio.create_task(server.serve(sockets=[sock]))
|
||||
try:
|
||||
async def ready():
|
||||
while not server.started:
|
||||
if task.done():
|
||||
await task
|
||||
raise RuntimeError("SERVER_START_FAILED")
|
||||
await asyncio.sleep(.01)
|
||||
await asyncio.wait_for(ready(), 10)
|
||||
report = await asyncio.wait_for(probe(f"http://127.0.0.1:{sock.getsockname()[1]}",
|
||||
credentials, size=size), 180)
|
||||
assert report["result"] == "PASSED"
|
||||
assert report["service_rss_bytes"] is None
|
||||
assert report["acceptance"] == "NOT_ASSESSED"
|
||||
assert len(report["transfers"]) == 4
|
||||
assert len({item["sha256"] for item in report["transfers"]}) == 4
|
||||
assert all(item["verified_bytes"] == size for item in report["transfers"])
|
||||
print(json.dumps(report))
|
||||
finally:
|
||||
server.should_exit = True
|
||||
await asyncio.wait_for(task, 10)
|
||||
try:
|
||||
asyncio.run(run())
|
||||
finally:
|
||||
sock.close()
|
||||
database.engine.dispose()
|
||||
@@ -0,0 +1,188 @@
|
||||
"""Four concurrent upload/download probes; run only against a disposable test service.
|
||||
|
||||
Credentials JSON is [{"username": "...", "password": "..."}, ...] for two
|
||||
existing test accounts. Never writes credentials, tokens or response bodies to reports.
|
||||
This is a transfer probe, not a claim of S-09 completion or an RSS measurement.
|
||||
"""
|
||||
import argparse
|
||||
import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
from pathlib import Path
|
||||
import time
|
||||
from urllib.parse import urlsplit
|
||||
import uuid
|
||||
|
||||
import httpx
|
||||
|
||||
CHUNK = 1024 * 1024
|
||||
|
||||
|
||||
class ProbeFailure(Exception):
|
||||
pass
|
||||
|
||||
|
||||
async def checked(client, method, path, **kwargs):
|
||||
response = await client.request(method, path, **kwargs)
|
||||
if response.status_code != 200:
|
||||
raise ProbeFailure(f"HTTP_{response.status_code}")
|
||||
return response.json()
|
||||
|
||||
|
||||
def block(index, offset, length):
|
||||
# Distinct, repeatable contents per stream and offset, without whole-file buffers.
|
||||
seed = hashlib.sha256(f"20260908:{index}:{offset}".encode()).digest()
|
||||
return (seed * ((length + len(seed) - 1) // len(seed)))[:length]
|
||||
|
||||
|
||||
def expected_hash(index, size):
|
||||
checksum = hashlib.sha256()
|
||||
for offset in range(0, size, CHUNK):
|
||||
checksum.update(block(index, offset, min(CHUNK, size - offset)))
|
||||
return checksum.hexdigest()
|
||||
|
||||
|
||||
async def probe(base_url, credentials, *, size=100 * CHUNK):
|
||||
if len(credentials) != 2 or credentials[0]["username"] == credentials[1]["username"]:
|
||||
raise ValueError("TWO_DISTINCT_TEST_ACCOUNTS_REQUIRED")
|
||||
if not 1 <= size <= 100 * CHUNK:
|
||||
raise ValueError("SIZE_OUT_OF_RANGE")
|
||||
run_id = uuid.uuid4().hex
|
||||
report = {"run_id": run_id, "seed": 20260908, "bytes_per_upload": size,
|
||||
"concurrency": 4, "accounts": 2, "vaults_per_account": 2,
|
||||
"result": "FAILED", "acceptance": "NOT_ASSESSED",
|
||||
"service_rss_bytes": None,
|
||||
"limitations": ["No server process/container RSS collector attached",
|
||||
"Not the 30-minute load, 10000-note or impaired-network benchmark",
|
||||
"Creates four test Vaults and retains uploaded data for inspection"],
|
||||
"transfers": []}
|
||||
timeout = httpx.Timeout(60, connect=10)
|
||||
async with httpx.AsyncClient(base_url=base_url, timeout=timeout,
|
||||
follow_redirects=False, trust_env=False) as admin:
|
||||
sessions = []
|
||||
for account in credentials:
|
||||
session = await checked(admin, "POST", "/sync/v1/auth/sessions", json={
|
||||
**account, "device_name": "upload-benchmark-" + run_id})
|
||||
sessions.append({"Authorization": "Bearer " + session["access_token"]})
|
||||
jobs = []
|
||||
for index in range(4):
|
||||
auth = sessions[index // 2]
|
||||
vault = await checked(admin, "POST", "/sync/v1/vaults", headers=auth,
|
||||
json={"name": f"benchmark-{run_id}-{index}"})
|
||||
jobs.append((index, auth, "/sync/v1/vaults/" + vault["vault_id"]))
|
||||
hashes = [expected_hash(index, size) for index in range(4)]
|
||||
gate = asyncio.Event()
|
||||
ready = asyncio.Queue()
|
||||
|
||||
async def transfer(index, auth, base):
|
||||
async with httpx.AsyncClient(base_url=base_url, headers=auth, timeout=timeout,
|
||||
follow_redirects=False, trust_env=False) as client:
|
||||
sha = hashes[index]
|
||||
info = await checked(client, "POST", base + "/uploads",
|
||||
json={"content_hash": sha, "size": size})
|
||||
if info["complete"] or info["offset"] != 0:
|
||||
raise ProbeFailure("NEW_VAULT_UPLOAD_NOT_EMPTY")
|
||||
upload = base + "/uploads/" + info["upload_id"]
|
||||
await ready.put(index)
|
||||
await gate.wait()
|
||||
start = time.perf_counter()
|
||||
for offset in range(0, size, CHUNK):
|
||||
data = block(index, offset, min(CHUNK, size - offset))
|
||||
result = await checked(client, "PUT", upload, params={"offset": offset}, content=data)
|
||||
if result["offset"] != offset + len(data):
|
||||
raise ProbeFailure("OFFSET_MISMATCH")
|
||||
receipt = await checked(client, "POST", upload + "/complete")
|
||||
if receipt != {"complete": True, "content_hash": sha}:
|
||||
raise ProbeFailure("COMPLETE_MISMATCH")
|
||||
if await checked(client, "POST", upload + "/complete") != receipt:
|
||||
raise ProbeFailure("COMPLETE_REPLAY_MISMATCH")
|
||||
uploaded = time.perf_counter()
|
||||
revision = {"operation_id": uuid.uuid4().hex, "file_id": uuid.uuid4().hex,
|
||||
"base_revision": 0, "path": "attachments/benchmark.bin",
|
||||
"operation": "put", "content_hash": sha, "size": size}
|
||||
committed = await checked(client, "POST", base + "/revisions", json=revision)
|
||||
if committed["sequence"] != 1:
|
||||
raise ProbeFailure("REVISION_SEQUENCE_MISMATCH")
|
||||
if await checked(client, "POST", base + "/revisions", json=revision) != committed:
|
||||
raise ProbeFailure("REVISION_REPLAY_MISMATCH")
|
||||
checksum, received = hashlib.sha256(), 0
|
||||
async with client.stream("GET", base + "/objects/" + sha) as response:
|
||||
if response.status_code != 200:
|
||||
raise ProbeFailure(f"DOWNLOAD_HTTP_{response.status_code}")
|
||||
async for data in response.aiter_bytes(CHUNK):
|
||||
received += len(data)
|
||||
if received > size:
|
||||
raise ProbeFailure("DOWNLOAD_SIZE_EXCEEDED")
|
||||
checksum.update(data)
|
||||
if received != size or checksum.hexdigest() != sha:
|
||||
raise ProbeFailure("DOWNLOAD_INTEGRITY")
|
||||
# Another account must not be able to read this object's bytes.
|
||||
denied = await admin.get(base + "/objects/" + sha, headers=sessions[1 - index // 2])
|
||||
if denied.status_code not in (403, 404):
|
||||
raise ProbeFailure("ACCOUNT_ISOLATION_FAILED")
|
||||
return {"index": index, "vault_id": base.rsplit("/", 1)[1],
|
||||
"sha256": sha, "verified_bytes": received,
|
||||
"upload_seconds": uploaded - start,
|
||||
"total_seconds": time.perf_counter() - start}
|
||||
|
||||
tasks = [asyncio.create_task(transfer(*job)) for job in jobs]
|
||||
try:
|
||||
# Bound preparation: a failed peer must not leave the others waiting forever.
|
||||
async def all_ready():
|
||||
for _ in jobs:
|
||||
await ready.get()
|
||||
readiness = asyncio.create_task(all_ready())
|
||||
try:
|
||||
finished, _ = await asyncio.wait([readiness, *tasks], timeout=75,
|
||||
return_when=asyncio.FIRST_COMPLETED)
|
||||
if not finished:
|
||||
raise ProbeFailure("PREPARATION_TIMEOUT")
|
||||
for task in finished:
|
||||
task.result()
|
||||
if not readiness.done():
|
||||
raise ProbeFailure("PREPARATION_FAILED")
|
||||
finally:
|
||||
readiness.cancel()
|
||||
await asyncio.gather(readiness, return_exceptions=True)
|
||||
start = time.perf_counter()
|
||||
gate.set()
|
||||
report["transfers"] = await asyncio.gather(*tasks)
|
||||
report["elapsed_seconds"] = time.perf_counter() - start
|
||||
report["result"] = "PASSED"
|
||||
return report
|
||||
finally:
|
||||
for task in tasks:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("--url", required=True)
|
||||
parser.add_argument("--credentials", type=Path, required=True)
|
||||
parser.add_argument("--output", type=Path, required=True)
|
||||
parser.add_argument("--allow-test-http", action="store_true")
|
||||
args = parser.parse_args()
|
||||
url = urlsplit(args.url)
|
||||
if (url.scheme not in ("http", "https") or not url.hostname or url.username or
|
||||
url.password or url.query or url.fragment or url.path not in ("", "/")):
|
||||
parser.error("EXPECTED_SERVICE_ORIGIN")
|
||||
if url.scheme == "http" and not args.allow_test_http:
|
||||
parser.error("HTTP_REQUIRES_ALLOW_TEST_HTTP")
|
||||
report = {"result": "FAILED", "acceptance": "NOT_ASSESSED"}
|
||||
try:
|
||||
credentials = json.loads(args.credentials.read_text(encoding="utf-8"))
|
||||
report = asyncio.run(asyncio.wait_for(probe(args.url, credentials), timeout=1800))
|
||||
except Exception as error:
|
||||
# Exception text may contain credentials or response data; log only its type.
|
||||
report["error_type"] = type(error).__name__
|
||||
finally:
|
||||
args.output.parent.mkdir(parents=True, exist_ok=True)
|
||||
args.output.write_text(json.dumps(report, indent=2) + "\n", encoding="utf-8")
|
||||
print(report["result"])
|
||||
return 0 if report["result"] == "PASSED" else 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
Reference in New Issue
Block a user