Files
imagefind/scripts/live-webdav-stress.py
T

280 lines
11 KiB
Python

#!/usr/bin/env python3
"""Stress ImageFind WebDAV reception without retaining generated large files."""
from __future__ import annotations
import argparse
import hashlib
import json
import threading
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
from typing import Any
from urllib.parse import quote
import httpx
BLOCK = bytes(1024 * 1024)
def padded_hash(prefix: bytes, size: int) -> str:
value = hashlib.sha256(prefix)
remaining = size - len(prefix)
while remaining > 0:
length = min(len(BLOCK), remaining)
value.update(BLOCK[:length])
remaining -= length
return value.hexdigest()
def padded_body(prefix: bytes, size: int):
yield prefix
remaining = size - len(prefix)
while remaining > 0:
length = min(len(BLOCK), remaining)
yield BLOCK[:length]
remaining -= length
def request(client: httpx.Client, method: str, path: str, **kwargs: Any) -> httpx.Response:
response = client.request(method, path, **kwargs)
if response.is_error:
raise RuntimeError(f"{method} {path} -> {response.status_code}: {response.text[:1000]}")
return response
def encoded(parts: list[str]) -> str:
return "/webdav/" + "/".join(quote(part, safe="") for part in parts)
def wait_uploads(api: httpx.Client, ids: set[str], timeout: int) -> dict[str, dict[str, Any]]:
deadline = time.monotonic() + timeout
terminal: dict[str, dict[str, Any]] = {}
while time.monotonic() < deadline:
rows = request(api, "GET", "/api/v1/uploads?limit=500").json()
by_id = {item["id"]: item for item in rows if item["id"] in ids}
terminal = {
key: value
for key, value in by_id.items()
if value["status"] in {"completed", "failed", "cancelled"}
}
if len(terminal) == len(ids):
failures = [item for item in terminal.values() if item["status"] != "completed"]
if failures:
raise AssertionError(f"committed WebDAV uploads failed: {failures}")
return terminal
time.sleep(1)
raise TimeoutError(f"WebDAV transfers did not finish: {ids - set(terminal)}")
def wait_video(api: httpx.Client, filename: str, timeout: int = 120) -> dict[str, Any]:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
rows = request(api, "GET", "/api/v1/videos?limit=500").json()
match = next((item for item in rows if item.get("source_key", "").endswith("/" + filename)), None)
if match:
return match
time.sleep(1)
raise TimeoutError(f"committed stress video did not appear: {filename}")
def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("--base-url", required=True)
parser.add_argument("--token-file", type=Path, required=True)
parser.add_argument("--run-dir", type=Path, required=True)
parser.add_argument("--case-id", required=True)
parser.add_argument("--size-mb", type=int, required=True)
parser.add_argument("--count", type=int, required=True)
parser.add_argument("--concurrency", type=int, default=4)
parser.add_argument("--commit-count", type=int, default=1)
parser.add_argument("--timeout", type=int, default=1800)
args = parser.parse_args()
state = json.loads((args.run_dir / "state.json").read_text())
token = args.token_file.read_text().strip()
prefix = f"{state['run_id']}-{args.case_id}-stress"
collection = f"{prefix}-collection"
group = "uploads"
size = args.size_mb * 1024**2
fixture = (args.run_dir / "e2e-positive.mp4").read_bytes()
if size < len(fixture):
raise ValueError("stress size must be at least the fixture size")
checksum = padded_hash(fixture, size)
auth = httpx.BasicAuth("imagefind", token)
api = httpx.Client(
base_url=args.base_url.rstrip("/"),
headers={"Authorization": f"Bearer {token}"},
timeout=httpx.Timeout(60, connect=10),
)
dav = httpx.Client(
base_url=args.base_url.rstrip("/"),
auth=auth,
timeout=httpx.Timeout(args.timeout, connect=15, read=args.timeout, write=args.timeout),
)
for parts in ([collection], [collection, group]):
response = dav.request("MKCOL", encoded(parts))
if response.status_code not in {201, 405}:
raise RuntimeError(f"MKCOL {parts} -> {response.status_code}: {response.text[:1000]}")
stop = threading.Event()
samples: list[dict[str, Any]] = []
def monitor() -> None:
with httpx.Client(
base_url=args.base_url.rstrip("/"),
headers={"Authorization": f"Bearer {token}"},
timeout=10,
) as client:
while not stop.wait(1):
started = time.perf_counter()
try:
status = client.get("/api/v1/status")
status.raise_for_status()
latency = round((time.perf_counter() - started) * 1000, 1)
resources = client.get("/api/v1/system/resources")
resources.raise_for_status()
payload = resources.json()
samples.append(
{
"api_ms": latency,
"cpu": payload.get("cpu_percent"),
"memory_available": payload.get("memory_available_bytes"),
"writer_wait_ms": (payload.get("database") or {}).get("last_wait_ms"),
"ok": True,
}
)
except Exception as exc:
samples.append({"ok": False, "error": type(exc).__name__})
monitor_thread = threading.Thread(target=monitor, name="imagefind-webdav-monitor", daemon=True)
monitor_thread.start()
files = [f"{prefix}-{index + 1:03d}.mp4" for index in range(args.count)]
temporary_paths = {name: encoded([collection, group, name + ".part"]) for name in files}
def put_one(filename: str) -> dict[str, Any]:
path = temporary_paths[filename]
started = time.perf_counter()
with httpx.Client(
base_url=args.base_url.rstrip("/"),
auth=auth,
timeout=httpx.Timeout(args.timeout, connect=15, read=args.timeout, write=args.timeout),
) as client:
response = client.put(
path,
content=padded_body(fixture, size),
headers={
"Content-Length": str(size),
"X-Content-SHA256": checksum,
},
)
if response.is_error:
raise RuntimeError(f"PUT {filename} -> {response.status_code}: {response.text[:1000]}")
if response.status_code not in {201, 204}:
raise RuntimeError(f"PUT {filename} returned {response.status_code}")
return {
"filename": filename,
"status": response.status_code,
"elapsed_ms": round((time.perf_counter() - started) * 1000, 1),
"offset": int(response.headers.get("upload-offset", 0)),
}
started_all = time.perf_counter()
results: list[dict[str, Any]] = []
error: Exception | None = None
committed: list[tuple[str, str]] = []
videos: list[dict[str, Any]] = []
try:
with ThreadPoolExecutor(max_workers=args.concurrency) as executor:
futures = {executor.submit(put_one, filename): filename for filename in files}
for future in as_completed(futures):
results.append(future.result())
if any(item["offset"] != size for item in results):
raise AssertionError("one or more WebDAV PUT responses reported an incomplete offset")
for filename in files[: min(args.commit_count, len(files))]:
source_path = temporary_paths[filename]
destination = encoded([collection, group, filename])
moved = request(
dav,
"MOVE",
source_path,
headers={
"Destination": args.base_url.rstrip("/") + destination,
"Overwrite": "F",
},
)
upload_id = moved.headers.get("x-imagefind-upload-id")
if not upload_id:
raise AssertionError(f"MOVE omitted upload id for {filename}")
committed.append((filename, upload_id))
if committed:
wait_uploads(api, {upload_id for _, upload_id in committed}, args.timeout)
for filename, _ in committed:
video = wait_video(api, filename)
ranged = request(
dav,
"GET",
encoded([collection, group, filename]),
headers={"Range": "bytes=0-63"},
)
if ranged.status_code != 206 or ranged.content != fixture[:64]:
raise AssertionError(f"committed range read failed: {filename}")
videos.append(video)
except Exception as exc:
error = exc
finally:
for filename in files[len(committed) :]:
try:
dav.delete(temporary_paths[filename])
except Exception:
pass
for video in videos:
try:
removed = request(api, "DELETE", f"/api/v1/videos/{video['id']}?delete_source=true").json()
if removed.get("trash_id"):
request(api, "DELETE", f"/api/v1/trash/{removed['trash_id']}")
except Exception:
pass
stop.set()
monitor_thread.join(timeout=15)
if error is not None:
raise error
elapsed = time.perf_counter() - started_all
successful_samples = [item for item in samples if item.get("ok")]
report = {
"case_id": args.case_id,
"size_mb": args.size_mb,
"count": args.count,
"concurrency": args.concurrency,
"committed": len(committed),
"elapsed_seconds": round(elapsed, 2),
"throughput_mib_s": round(args.size_mb * args.count / max(elapsed, 0.001), 2),
"put_latency_ms": {
"minimum": min(item["elapsed_ms"] for item in results),
"maximum": max(item["elapsed_ms"] for item in results),
},
"monitor": {
"samples": len(samples),
"failed_samples": sum(1 for item in samples if not item.get("ok")),
"max_api_ms": max((item["api_ms"] for item in successful_samples), default=0),
"max_cpu_percent": max((item["cpu"] for item in successful_samples), default=0),
"min_memory_available_gb": round(
min((item["memory_available"] for item in successful_samples), default=0) / 1024**3,
2,
),
"max_writer_wait_ms": max((item["writer_wait_ms"] for item in successful_samples), default=0),
},
}
report_path = args.run_dir / f"webdav-stress-{args.case_id}.json"
report_path.write_text(json.dumps(report, ensure_ascii=False, indent=2) + "\n")
print(json.dumps(report, ensure_ascii=False, indent=2))
if __name__ == "__main__":
main()