#!/usr/bin/env python3 """Stage, upload, and verify the public GLM SQG calibration dataset. ``launch`` hardlinks the sealed raw capture/reproduction inputs into one local dataset tree and starts a detached resumable upload. ``finish`` waits for that initial transfer, hardlinks the derived evidence, writes a complete hashed local census, resumes the upload, and verifies the remote revision census and commit. The Hub token is read only from a root-readable file and is never a command-line argument or serialized artifact. """ from __future__ import annotations import argparse from concurrent.futures import ThreadPoolExecutor import hashlib import json import os from pathlib import Path import shutil import subprocess import sys import textwrap import time from typing import Any, Iterable, Mapping if __package__: from b300_remote.hf_cli import resolve_hf_cli else: # Direct paid-node script invocation. from hf_cli import resolve_hf_cli SCHEMA = "glm52-bmm-law-sqg-hessian-dataset-v1" SOURCE_REVISION = "b4734de4facf877f85769a911abafc5283eab3d9" def atomic_json(path: Path, value: object) -> None: path.parent.mkdir(parents=True, exist_ok=True) temporary = path.with_name(f".{path.name}.tmp-{os.getpid()}") with temporary.open("x", encoding="utf-8") as handle: json.dump(value, handle, indent=2, sort_keys=True) handle.write("\n") handle.flush() os.fsync(handle.fileno()) os.replace(temporary, path) def atomic_text(path: Path, value: str) -> None: path.parent.mkdir(parents=True, exist_ok=True) temporary = path.with_name(f".{path.name}.tmp-{os.getpid()}") with temporary.open("x", encoding="utf-8") as handle: handle.write(value) handle.flush() os.fsync(handle.fileno()) os.replace(temporary, path) def write_dataset_card( root: Path, repo: str, revision: str, source_revision: str, *, accepted_tag: str | None = None, ) -> None: atomic_text( root / "README.md", textwrap.dedent( f"""\ --- license: other license_name: glm-5.2-derived-calibration-data license_link: https://huggingface.co/zai-org/GLM-5.2/blob/{source_revision}/LICENSE task_categories: [text-generation] tags: [glm, calibration, hessian, sqg, bmm-law, quantization] pretty_name: GLM-5.2 BMM Law SQG Hessians --- # GLM-5.2 BMM Law SQG Hessians Public calibration and reconstruction evidence for the all-SQG GLM-5.2 W4A8 quantization campaign. It includes the sealed PP8/TP1 routed capture, MTP layer-78 capture, dense Hessians, document split plan, derived H13/H2/(H,B) evidence, exact profile and allocation records, and reproduction sources. - Dataset repository: `{repo}` - Immutable build revision: `{revision}` - Public default revision: `main` - Immutable accepted tag: `{accepted_tag or "created after final verification"}` - Official BF16 source: `zai-org/GLM-5.2@{source_revision}` - Capture topology: TP1/PP8 - Quantization target: all-SQG full W4A8; routed K3/K4 and dense K6 Large capture objects are identified by their verified Hub object identities at the final immutable commit. Smaller evidence files additionally carry local SHA-256 values in `HESSIAN_DATASET_MANIFEST.json`. No Hugging Face token is stored in this dataset. ## License and third-party rights This publication does not relicense GLM-5.2 or bundled third-party reproduction code. Model-derived calibration data remains subject to the upstream GLM-5.2 terms and applicable third-party notices. See `LICENSE` and notices within reproduced source trees. """ ), ) atomic_text( root / "LICENSE", textwrap.dedent( f"""\ GLM-5.2-DERIVED CALIBRATION DATA NOTICE This dataset contains calibration artifacts derived from zai-org/GLM-5.2 at revision {source_revision} and reproduction code from multiple upstream projects. No new license is asserted over the upstream model or third-party code. Use, redistribution, and modification remain subject to the upstream GLM-5.2 license and the notices shipped with each source: https://huggingface.co/zai-org/GLM-5.2/blob/{source_revision}/LICENSE """ ), ) def parse_source(raw: str) -> tuple[str, Path]: if "=" not in raw: raise ValueError("source must be REPO_PREFIX=/absolute/path") prefix, path = raw.split("=", 1) if ( not prefix or prefix.startswith("/") or ".." in Path(prefix).parts or not Path(path).is_absolute() ): raise ValueError(f"unsafe dataset source: {raw}") return prefix.strip("/"), Path(path).resolve() def link_file(source: Path, destination: Path) -> None: if source.is_symlink() or not source.is_file(): raise ValueError(f"dataset source is not a regular file: {source}") destination.parent.mkdir(parents=True, exist_ok=True) if destination.exists(): left = source.stat() right = destination.stat() if (left.st_dev, left.st_ino, left.st_size) != ( right.st_dev, right.st_ino, right.st_size, ): raise ValueError(f"existing dataset hardlink differs: {destination}") return try: os.link(source, destination, follow_symlinks=False) except OSError as exc: raise OSError( f"dataset/source must share one filesystem for zero-copy staging: {source}" ) from exc def stage_sources(root: Path, sources: Iterable[tuple[str, Path]]) -> int: count = 0 for prefix, source in sources: if source.is_file(): link_file(source, root / prefix / source.name) count += 1 continue if not source.is_dir() or source.is_symlink(): raise FileNotFoundError(source) for path in sorted(source.rglob("*")): relative = path.relative_to(source) if "__pycache__" in relative.parts or path.suffix in {".pyc", ".pyo"}: continue if path.is_symlink(): raise ValueError(f"dataset source tree contains symlink: {path}") if path.is_file(): link_file(path, root / prefix / relative) count += 1 return count def token_environment(token_file: Path) -> dict[str, str]: if token_file.is_symlink() or not token_file.is_file(): raise ValueError("HF token file must be a non-symlink with mode 0600/0400") token_file = token_file.resolve(strict=True) stat = token_file.stat() if stat.st_mode & 0o077: raise ValueError("HF token file must be a non-symlink with mode 0600/0400") token = token_file.read_text(encoding="utf-8").strip() if not token or any(char.isspace() for char in token): raise ValueError("HF token file is empty or malformed") env = os.environ.copy() env["HF_TOKEN"] = token env["HF_XET_HIGH_PERFORMANCE"] = "1" return env def upload(root: Path, repo: str, revision: str, token_file: Path, workers: int) -> None: hf_cli = resolve_hf_cli() env = token_environment(token_file) # Public is an enforced remote property, not merely a claim in the local # receipt. Keep the token in the subprocess environment and out of argv. public_script = ( "from huggingface_hub import HfApi; " f"a=HfApi(); a.create_repo({repo!r}, repo_type='dataset', " "private=False, exist_ok=True); " f"a.update_repo_settings({repo!r}, repo_type='dataset', private=False)" ) subprocess.run([sys.executable, "-c", public_script], check=True, env=env) subprocess.run( [hf_cli, "repos", "branch", "create", repo, revision, "--type", "dataset", "--exist-ok"], check=True, env=env, ) subprocess.run( [ hf_cli, "upload-large-folder", repo, str(root), "--type", "dataset", "--revision", revision, "--num-workers", str(workers), "--no-bars", ], check=True, env=env, ) def worker(args: argparse.Namespace) -> int: result = args.state_root / "initial-upload-result.json" try: upload(args.dataset_root, args.repo, args.revision, args.token_file, args.workers) local = census(args.dataset_root, args.workers) remote = verify_remote( args.dataset_root, args.repo, args.revision, args.token_file, local, ) value = { "complete": True, "initial_upload_verified": True, "repo": args.repo, "revision": args.revision, "commit": remote["commit"], "remote_object_domain_sha256": remote["object_domain_sha256"], "file_count": remote["file_count"], "total_bytes": remote["total_bytes"], "public": True, "token_serialized": False, "uploader_source_sha256": sha256_file(Path(__file__).resolve()), } except BaseException as exc: value = { "complete": False, "initial_upload_verified": False, "repo": args.repo, "revision": args.revision, "public": False, "error": repr(exc), "token_serialized": False, "uploader_source_sha256": sha256_file(Path(__file__).resolve()), } atomic_json(result, value) raise atomic_json(result, value) return 0 def sha256_file(path: Path) -> str: digest = hashlib.sha256() with path.open("rb") as handle: while block := handle.read(64 << 20): digest.update(block) return digest.hexdigest() def census( root: Path, workers: int, *, hash_max_bytes: int = 64 << 20, ) -> list[dict[str, object]]: # `hf upload-large-folder` owns a resumable local cache beneath # `.cache/huggingface`. It is deliberately not uploaded, so including it # in the sealed dataset census would make remote verification fail after a # successful transfer. Exclude only that exact tool-owned subtree; every # actual dataset file remains mandatory. paths = sorted( path for path in root.rglob("*") if path.is_file() and path.relative_to(root).parts[:2] != (".cache", "huggingface") ) def one(path: Path) -> dict[str, object]: size = path.stat().st_size return { "path": path.relative_to(root).as_posix(), "bytes": size, # Avoid a redundant full-terabyte local read. The Hub uploader # content-hashes large objects, and verify_remote binds those # LFS/Xet identities to the immutable final commit. "sha256": sha256_file(path) if size <= hash_max_bytes else None, "identity": ( "local_sha256" if size <= hash_max_bytes else "verified_hub_object_at_final_commit" ), } with ThreadPoolExecutor(max_workers=workers) as pool: return list(pool.map(one, paths)) def launch(args: argparse.Namespace) -> int: # Fail before staging the large dataset or publishing a launch receipt. resolve_hf_cli() token_environment(args.token_file) args.dataset_root.mkdir(parents=True, exist_ok=True) args.state_root.mkdir(parents=True, exist_ok=True) write_dataset_card(args.dataset_root, args.repo, args.revision, args.source_revision) count = stage_sources(args.dataset_root, map(parse_source, args.source)) result = args.state_root / "initial-upload-result.json" prior_launch: dict[str, object] = {} if args.receipt.is_file(): candidate = json.loads(args.receipt.read_text(encoding="utf-8")) if ( candidate.get("repo") == args.repo and candidate.get("revision") == args.revision ): prior_launch = candidate def process_alive(value: object) -> bool: if not isinstance(value, int) or value <= 1: return False try: os.kill(value, 0) except (OSError, ValueError): return False try: command = (Path("/proc") / str(value) / "cmdline").read_bytes().split(b"\0") except OSError: return False decoded = [item.decode(errors="replace") for item in command if item] return ( any(Path(item).name == Path(__file__).name for item in decoded) and "worker" in decoded and args.repo in decoded and args.revision in decoded ) if result.is_file() and json.loads(result.read_text()).get("complete") is True: pid = None elif process_alive(prior_launch.get("pid")): # A supervisor retry must adopt, never duplicate, an existing detached # upload worker. Its resumable cache belongs to that single process. pid = int(prior_launch["pid"]) else: log = (args.state_root / "initial-upload.log").open("ab", buffering=0) command = [ sys.executable, str(Path(__file__).resolve()), "worker", "--dataset-root", str(args.dataset_root), "--state-root", str(args.state_root), "--repo", args.repo, "--revision", args.revision, "--token-file", str(args.token_file), "--workers", str(args.workers), ] process = subprocess.Popen( command, stdin=subprocess.DEVNULL, stdout=log, stderr=subprocess.STDOUT, start_new_session=True, close_fds=True, ) pid = process.pid initial_verified = False initial_commit = None if result.is_file(): prior = json.loads(result.read_text(encoding="utf-8")) initial_verified = prior.get("initial_upload_verified") is True initial_commit = prior.get("commit") if initial_verified else None atomic_json( args.receipt, { "schema": SCHEMA, "complete": True, "phase": "initial_upload_launch", "launch_started": True, "initial_upload_verified": initial_verified, "repo": args.repo, "revision": args.revision, "commit": initial_commit, # A detached process launch is not evidence that anything is public. "public": True if initial_verified else False, "staged_file_count": count, "pid": pid, "token_serialized": False, "uploader_source_sha256": sha256_file(Path(__file__).resolve()), }, ) return 0 def _remote_snapshot( repo: str, revision: str, token_file: Path ) -> dict[str, Any]: """Read one exact public dataset revision without serializing credentials.""" env = token_environment(token_file) script = ( "import json; from huggingface_hub import HfApi; " "i=HfApi().dataset_info(" + repr(repo) + ", revision=" + repr(revision) + ", files_metadata=True); " "print(json.dumps({'sha':i.sha,'private':i.private,'files':{s.rfilename:" "{'size':s.size,'blob_id':s.blob_id,'lfs_sha256':getattr(s.lfs,'sha256',None)} " "for s in i.siblings}}))" ) result = subprocess.run( [sys.executable, "-c", script], check=True, capture_output=True, text=True, env=env, ) value = json.loads(result.stdout) if value.get("private") is not False: raise ValueError("Hessian dataset repository is not public") return value def _object_domain(files: Mapping[str, Mapping[str, object]]) -> dict[str, Any]: normalized = { str(path): { "size": int(item["size"]), "blob_id": item.get("blob_id"), "lfs_sha256": item.get("lfs_sha256"), } for path, item in sorted(files.items()) } return { "files": normalized, "object_domain_sha256": hashlib.sha256( json.dumps(normalized, sort_keys=True, separators=(",", ":")).encode() ).hexdigest(), "file_count": len(normalized), "total_bytes": sum(int(item["size"]) for item in normalized.values()), } def verify_remote( root: Path, repo: str, revision: str, token_file: Path, local: list[dict[str, object]] ) -> dict[str, Any]: del root # Retained in the public signature for existing callers. remote = _remote_snapshot(repo, revision, token_file) expected = {str(item["path"]): int(item["bytes"]) for item in local} observed = { str(key): int(value["size"]) for key, value in remote["files"].items() } # Hub-generated metadata is allowed; every local dataset file is mandatory. missing = {key: size for key, size in expected.items() if observed.get(key) != size} if missing: raise ValueError(f"remote Hessian dataset census differs for {len(missing)} files") extras = set(observed) - set(expected) if extras - {".gitattributes"}: raise ValueError(f"remote Hessian dataset has unexpected files: {sorted(extras)[:8]}") domain = _object_domain(remote["files"]) domain["commit"] = str(remote["sha"]) domain["payload_file_count"] = len(expected) domain["payload_total_bytes"] = sum(expected.values()) return domain def _assert_initial_revision_unchanged( *, repo: str, revision: str, token_file: Path, status: Mapping[str, object] ) -> None: del revision commit = status.get("commit") if not isinstance(commit, str) or len(commit) != 40: raise ValueError("initial upload verification lacks an immutable commit") # Verify the immutable commit itself. The mutable build branch may already # have advanced during a previously interrupted final sync. current = _remote_snapshot(repo, commit, token_file) domain = _object_domain(current["files"]) if ( current.get("sha") != status.get("commit") or domain["object_domain_sha256"] != status.get("remote_object_domain_sha256") ): raise ValueError("initial public Hessian build revision changed after verification") def _verify_or_repair_initial_upload(args: argparse.Namespace) -> dict[str, Any]: """Adopt a verified launch, or resume a legacy/unverified launch safely.""" result = args.state_root / "initial-upload-result.json" if result.is_file(): status = json.loads(result.read_text(encoding="utf-8")) if status.get("initial_upload_verified") is True: _assert_initial_revision_unchanged( repo=args.repo, revision=args.revision, token_file=args.token_file, status=status, ) return status if status.get("complete") is False: # A failed detached upload is resumable; do not make the old failure # permanent and do not discard its upload-large-folder cache. pass local = census(args.dataset_root, args.hash_workers, hash_max_bytes=args.hash_max_bytes) upload(args.dataset_root, args.repo, args.revision, args.token_file, args.workers) remote = verify_remote( args.dataset_root, args.repo, args.revision, args.token_file, local ) status = { "complete": True, "initial_upload_verified": True, "repo": args.repo, "revision": args.revision, "commit": remote["commit"], "remote_object_domain_sha256": remote["object_domain_sha256"], "file_count": remote["file_count"], "total_bytes": remote["total_bytes"], "public": True, "token_serialized": False, "uploader_source_sha256": sha256_file(Path(__file__).resolve()), } atomic_json(result, status) return status def publish_main_and_tag( *, root: Path, repo: str, build_revision: str, build: Mapping[str, Any], accepted_tag: str, token_file: Path, workers: int, ) -> dict[str, Any]: """Publish the verified build payload to main through the resumable uploader. ``CommitOperationCopy`` resolves every source path separately. A complete GLM Hessian dataset has more files than the Hub's resolver-rate window, so a one-commit copy can never finish. The large-folder uploader uses the already uploaded object identities and deduplicates the payload without issuing one resolver request per file. """ if not build_revision.startswith("build-"): raise ValueError("Hessian publication source must be a build revision") token = token_environment(token_file)["HF_TOKEN"] from huggingface_hub import CommitOperationDelete, HfApi api = HfApi(token=token) main_info = api.dataset_info(repo, revision="main", files_metadata=True) main_files = { item.rfilename: { "size": item.size, "blob_id": item.blob_id, "lfs_sha256": getattr(item.lfs, "sha256", None), } for item in main_info.siblings } desired = dict(build["files"]) if _object_domain(main_files)["object_domain_sha256"] != build["object_domain_sha256"]: # `upload-large-folder` stores commit state below the uploaded root and # does not namespace that state by revision. Reusing the build root for # `main` therefore reports every build object as already committed and # leaves main unchanged. A hardlink-only sibling gives main an # independent cache without duplicating the multi-terabyte payload. main_root = root.with_name(f"{root.name}-main-publication") if main_root.is_symlink() or (main_root.exists() and not main_root.is_dir()): raise ValueError(f"unsafe main publication staging root: {main_root}") if main_root.exists(): shutil.rmtree(main_root) shutil.copytree( root, main_root, copy_function=os.link, ignore=shutil.ignore_patterns(".cache"), ) upload(main_root, repo, "main", token_file, workers) refreshed = _remote_snapshot(repo, "main", token_file) stale = sorted(set(refreshed["files"]) - set(desired)) if stale: api.create_commit( repo, repo_type="dataset", revision="main", parent_commit=str(refreshed["sha"]), operations=[ CommitOperationDelete(path_in_repo=path) for path in stale ], commit_message="Remove stale files from accepted Hessian dataset", ) main = _remote_snapshot(repo, "main", token_file) main_domain = _object_domain(main["files"]) if main_domain["object_domain_sha256"] != build["object_domain_sha256"]: raise ValueError("public main Hessian object domain differs from verified build") api.create_tag( repo, tag=accepted_tag, revision=str(main["sha"]), repo_type="dataset", exist_ok=True, ) tagged = api.dataset_info(repo, revision=accepted_tag, files_metadata=True) if tagged.sha != main["sha"]: raise ValueError("accepted Hessian dataset tag does not resolve to main commit") main_domain["commit"] = str(main["sha"]) main_domain["accepted_tag"] = accepted_tag return main_domain def finish(args: argparse.Namespace) -> int: accepted_tag = args.accepted_tag if accepted_tag is None: if not args.revision.startswith("build-") or len(args.revision) <= len("build-"): raise ValueError("cannot derive accepted tag from non-build revision") accepted_tag = "accepted-" + args.revision.removeprefix("build-") initial = args.state_root / "initial-upload-result.json" deadline = time.monotonic() + args.initial_wait_seconds while not initial.is_file(): if time.monotonic() >= deadline: raise TimeoutError("initial Hessian upload did not publish a result before timeout") time.sleep(10) # Legacy launch receipts said complete when only Popen had succeeded. This # gate either verifies the exact public branch or safely resumes it. status = _verify_or_repair_initial_upload(args) if status.get("initial_upload_verified") is not True: raise RuntimeError("initial Hessian upload is not remotely verified") write_dataset_card( args.dataset_root, args.repo, args.revision, args.source_revision, accepted_tag=accepted_tag, ) stage_sources(args.dataset_root, map(parse_source, args.source)) local = [ item for item in census( args.dataset_root, args.hash_workers, hash_max_bytes=args.hash_max_bytes, ) if item["path"] != "HESSIAN_DATASET_MANIFEST.json" ] manifest = { "schema": SCHEMA, "complete": True, "repo": args.repo, "revision": args.revision, "public": True, "official_bf16_revision": args.source_revision, "file_count": len(local), "total_bytes": sum(int(item["bytes"]) for item in local), "files": local, "large_payload_identity": "verified_hub_object_at_final_commit", "contains": [ "raw_capture", "derived_h13", "derived_h2", "derived_cross_term_h_b", "split_plan", "source_revision", "schemas_hashes", "reproduction_scripts", ], "uploader_source_sha256": sha256_file(Path(__file__).resolve()), } atomic_json(args.dataset_root / "HESSIAN_DATASET_MANIFEST.json", manifest) # Include the manifest itself in the final local/remote census. local = census( args.dataset_root, args.hash_workers, hash_max_bytes=args.hash_max_bytes, ) upload(args.dataset_root, args.repo, args.revision, args.token_file, args.workers) build = verify_remote( args.dataset_root, args.repo, args.revision, args.token_file, local ) main = publish_main_and_tag( root=args.dataset_root, repo=args.repo, build_revision=args.revision, build=build, accepted_tag=accepted_tag, token_file=args.token_file, workers=args.workers, ) atomic_json( args.receipt, { "schema": SCHEMA, "complete": True, "phase": "final_publication_verified", "launch_started": True, "initial_upload_verified": True, "repo": args.repo, "revision": "main", "commit": main["commit"], "build_revision": args.revision, "build_commit": build["commit"], "build_verified": True, "main_commit": main["commit"], "main_verified": True, "accepted_tag": accepted_tag, "accepted_tag_verified": True, "build_remote_object_domain_sha256": build["object_domain_sha256"], "main_remote_object_domain_sha256": main["object_domain_sha256"], "remote_object_domain_sha256": main["object_domain_sha256"], "public": True, "file_count": main["file_count"], "total_bytes": main["total_bytes"], "payload_file_count": len(local), "payload_total_bytes": sum(int(item["bytes"]) for item in local), "manifest_sha256": sha256_file( args.dataset_root / "HESSIAN_DATASET_MANIFEST.json" ), "uploader_source_sha256": sha256_file(Path(__file__).resolve()), "token_serialized": False, }, ) return 0 def parser() -> argparse.ArgumentParser: result = argparse.ArgumentParser(description=__doc__) sub = result.add_subparsers(dest="command", required=True) for name in ("launch", "finish"): command = sub.add_parser(name) command.add_argument("--dataset-root", type=Path, required=True) command.add_argument("--state-root", type=Path, required=True) command.add_argument("--repo", default="brandonmusic/GLM-5.2-BMM-Law-SQG-Hessians") command.add_argument("--revision", required=True) command.add_argument("--token-file", type=Path, default=Path("/run/secrets/glm52-hf-token")) command.add_argument("--workers", type=int, default=32) command.add_argument("--source", action="append", default=[]) command.add_argument("--receipt", type=Path, required=True) command.add_argument("--source-revision", default=SOURCE_REVISION) finish_parser = sub.choices["finish"] finish_parser.add_argument("--hash-workers", type=int, default=32) finish_parser.add_argument("--hash-max-bytes", type=int, default=64 << 20) finish_parser.add_argument("--initial-wait-seconds", type=float, default=7_200) # Backward-compatible with an already-running supervisor whose immutable # build revision was sealed before this safety correction. finish_parser.add_argument("--accepted-tag") worker_parser = sub.add_parser("worker") worker_parser.add_argument("--dataset-root", type=Path, required=True) worker_parser.add_argument("--state-root", type=Path, required=True) worker_parser.add_argument("--repo", required=True) worker_parser.add_argument("--revision", required=True) worker_parser.add_argument("--token-file", type=Path, required=True) worker_parser.add_argument("--workers", type=int, required=True) return result def main() -> int: args = parser().parse_args() for path_name in ("dataset_root", "state_root", "receipt"): if hasattr(args, path_name): setattr(args, path_name, getattr(args, path_name).resolve()) if args.command == "launch": return launch(args) if args.command == "finish": return finish(args) return worker(args) if __name__ == "__main__": raise SystemExit(main())