#!/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()