# /// script # requires-python = ">=3.12" # dependencies = ["httpx>=0.28,<1", "psutil>=6,<8"] # /// """Run: uv run report.py --mode numeric (or --mode checklist). Replace the simulated work with your actual records/files/subtasks. State contains a write token; keep it private, outside version control. A restart resumes the same page. """ import argparse import json import os import time from datetime import UTC, datetime from pathlib import Path import httpx import psutil def metrics(disk_path: str) -> dict: memory = psutil.virtual_memory() disk = psutil.disk_usage(disk_path) return { "cpu_percent": psutil.cpu_percent(interval=None), "memory_percent": memory.percent, "memory_used_bytes": memory.used, "memory_total_bytes": memory.total, "disk_percent": disk.percent, "disk_used_bytes": disk.used, "disk_total_bytes": disk.total, "measured_at": datetime.now(UTC).isoformat(), } def request(client, method, path, *, payload=None, token=None, creating=False): headers = {"Authorization": f"Bearer {token}"} if token else {} for attempt in range(6): try: response = client.request(method, path, json=payload, headers=headers) except httpx.TransportError: # Creation may have succeeded: blindly retrying could create a page whose token was lost. if creating or attempt == 5: raise time.sleep(min(2**attempt, 15)) continue if response.status_code == 429: time.sleep(float(response.headers.get("Retry-After", "1")) + 0.05) continue if response.status_code >= 500 and not creating and attempt < 5: time.sleep(min(2**attempt, 15)) continue response.raise_for_status() return response.json() raise RuntimeError("Rate limit retry budget exhausted; retry the job later") def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--base-url", default="https://progress.parklab.work") parser.add_argument("--state", type=Path, default=Path(".progress-state.json")) parser.add_argument("--mode", choices=["numeric", "checklist"], default="numeric") parser.add_argument("--total", type=int, default=10000) parser.add_argument("--delay", type=float, default=0.002) parser.add_argument("--disk-path", default=str(Path.cwd().anchor)) args = parser.parse_args() if args.total <= 0 or args.delay < 0: parser.error("total must be positive and delay must be nonnegative") with httpx.Client(base_url=args.base_url, timeout=15, follow_redirects=False) as client: if args.state.exists(): state = json.loads(args.state.read_text()) if state["base_url"] != args.base_url: raise ValueError("State belongs to a different API server") current = request(client, "GET", f"/api/v1/progress/{state['id']}") else: payload = {"title": "Python task progress", "precision": 6, "mode": args.mode} if args.mode == "numeric": payload.update(total=args.total, unit="records") else: payload["items"] = [ {"title": "Validate input", "weight": 1}, {"title": "Process records", "total": args.total, "weight": 8}, {"title": "Save results", "weight": 1}, ] created = request(client, "POST", "/api/v1/progress", payload=payload, creating=True) current = created["progress"] state = { "id": current["id"], "write_token": created["write_token"], "url": created["url"], "base_url": args.base_url, } fd = os.open(args.state, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) with os.fdopen(fd, "w") as file: json.dump(state, file) print("Progress:", state["url"]) print("CSV:", f"{args.base_url}/api/v1/progress/{state['id']}/history.csv") prefix = f"/api/v1/progress/{state['id']}" token = state["write_token"] psutil.cpu_percent(interval=None) # Prime the first non-blocking CPU sample. last_call = 0.0 def write(path, payload, method="POST"): nonlocal last_call # Leave headroom below the shared 5/s IP cap; many reporters on one IP must coordinate. time.sleep(max(0, 0.3 - (time.monotonic() - last_call))) try: return request(client, method, prefix + path, payload=payload, token=token) finally: last_call = time.monotonic() try: write("/reports", {"status": "running", "message": "Starting work"}) if current["mode"] == "checklist": first, work, last = current["items"] write(f"/items/{first['id']}", {"status": "completed"}, "PATCH") start = int(float(work["completed"])) total = int(float(work["total"])) if work["status"] != "completed": write(f"/items/{work['id']}", {"status": "running"}, "PATCH") else: start, total = int(float(current["completed"])), int(float(current["total"])) last_report = time.monotonic() for completed in range(start + 1, total + 1): time.sleep( args.delay ) # Replace with process_record(...), not simulated percentages. if time.monotonic() - last_report < 1 and completed != total: continue if current["mode"] == "checklist": write( f"/items/{work['id']}", {"completed": completed, "message": f"{completed:,} records processed"}, "PATCH", ) write("/reports", {"metrics": metrics(args.disk_path)}) else: write( "/reports", { "completed": completed, "metrics": metrics(args.disk_path), "message": f"{completed:,} records processed", }, ) last_report = time.monotonic() if current["mode"] == "checklist": write(f"/items/{work['id']}", {"status": "completed"}, "PATCH") write(f"/items/{last['id']}", {"status": "completed"}, "PATCH") write( "/reports", { "status": "completed", "message": "All work completed", "metrics": metrics(args.disk_path), }, ) except Exception: try: write("/reports", {"status": "failed", "message": "Work failed: check reporter logs"}) except (httpx.HTTPError, RuntimeError): print("Failure report could not be delivered; original exception follows.") raise if __name__ == "__main__": main()