review round: pointer written+read-back before prune, verify retried, latest rolled to final, stale-pointer scan on resume, ddp_timeout 4h, chunked batched commits, manifest self-consistency gate, full content verification after upload
Browse files- build/publish_mix.py +129 -51
build/publish_mix.py
CHANGED
|
@@ -120,21 +120,29 @@ def attribution(man):
|
|
| 120 |
def remote_digests(api, repo):
|
| 121 |
"""path -> content sha256 for what is already on the Hub, so a re-run resumes instead of resending.
|
| 122 |
|
| 123 |
-
Only LFS-tracked files expose a content sha256 (`lfs.oid`, sha256
|
| 124 |
-
files expose only a git blob sha, which is not a content sha256, so they
|
| 125 |
-
|
| 126 |
-
huggingface_hub, so attribute access alone silently returns None for every shard and
|
| 127 |
-
into a full 2.3 GB re-upload (Gate 2f).
|
|
|
|
|
|
|
|
|
|
|
|
|
| 128 |
"""
|
| 129 |
out = {}
|
| 130 |
for ent in api.dataset_info(repo, repo_type="dataset", files_metadata=True).siblings or []:
|
| 131 |
lfs = getattr(ent, "lfs", None)
|
| 132 |
if isinstance(lfs, dict):
|
| 133 |
-
|
| 134 |
elif lfs is not None:
|
| 135 |
-
|
| 136 |
else:
|
| 137 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 138 |
return out
|
| 139 |
|
| 140 |
|
|
@@ -187,6 +195,18 @@ def main():
|
|
| 187 |
f"{str(aud.get('mix_bytes_sha256'))[:16]}, this manifest is {want_digest[:16]}")
|
| 188 |
if int(aud.get("val_documents_audited") or 0) <= 0:
|
| 189 |
problems.append("the audit never scanned the held-out validation shards")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 190 |
if problems:
|
| 191 |
print("FATAL: contamination gate refuses publication:", flush=True)
|
| 192 |
for p in problems:
|
|
@@ -262,35 +282,43 @@ def main():
|
|
| 262 |
print(f" (no resumable listing: {type(e).__name__}; uploading everything)", flush=True)
|
| 263 |
|
| 264 |
todo = []
|
|
|
|
| 265 |
for np, lp in to_upload:
|
| 266 |
h = hashlib.sha256()
|
| 267 |
with open(lp, "rb") as fh:
|
| 268 |
for chunk in iter(lambda: fh.read(1 << 20), b""):
|
| 269 |
h.update(chunk)
|
| 270 |
-
local_sha = h.hexdigest()
|
| 271 |
if remote_sha.get(np) == local_sha:
|
| 272 |
print(f" = {np} already on the Hub with this sha256, skipped", flush=True)
|
| 273 |
continue
|
| 274 |
todo.append((np, lp, local_sha))
|
| 275 |
print(f" {len(todo)} of {len(to_upload)} files still need uploading", flush=True)
|
| 276 |
|
| 277 |
-
#
|
| 278 |
-
# `upload_file` commits succeeded and then the Hub
|
| 279 |
-
# files -- with a bare HfHubHTTPError,
|
| 280 |
-
#
|
|
|
|
| 281 |
failures = []
|
| 282 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 283 |
try:
|
| 284 |
-
api.upload_folder(repo_id=a.repo, repo_type="dataset", folder_path=a.root,
|
| 285 |
-
|
| 286 |
-
|
| 287 |
-
|
| 288 |
-
|
| 289 |
-
|
|
|
|
|
|
|
| 290 |
except Exception as e:
|
| 291 |
print(f" batched commit failed ({type(e).__name__}: {str(e)[:160]}); "
|
| 292 |
"falling back to per-file uploads with retries", flush=True)
|
| 293 |
-
for np, lp, local_sha in
|
| 294 |
for attempt in (1, 2, 3):
|
| 295 |
try:
|
| 296 |
with open(lp, "rb") as fh:
|
|
@@ -313,33 +341,66 @@ def main():
|
|
| 313 |
sys.exit(11)
|
| 314 |
|
| 315 |
# Independent verification from the Hub, not from the client that just wrote it. A batched commit can
|
| 316 |
-
# return without raising and still not carry every file, so presence is checked per path and
|
| 317 |
-
#
|
| 318 |
-
|
| 319 |
-
|
| 320 |
-
|
| 321 |
-
|
| 322 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 323 |
want = {np for np, _ in to_upload} | {".gitattributes"}
|
| 324 |
local_by_path = {np: lp for np, lp in to_upload}
|
| 325 |
-
|
| 326 |
-
|
| 327 |
-
|
| 328 |
-
|
| 329 |
-
|
| 330 |
-
|
| 331 |
-
|
| 332 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 333 |
sha_bad.append(np)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 334 |
R = {"repo": a.repo, "files_expected": len(want), "files_present": len(want & remote),
|
| 335 |
-
"missing":
|
| 336 |
-
"
|
| 337 |
-
|
| 338 |
-
|
|
|
|
|
|
|
| 339 |
first = man["shards"][0]["file"]
|
| 340 |
-
url = f"https://huggingface.co/datasets/{a.repo}/resolve/main/{first}"
|
| 341 |
try:
|
| 342 |
-
with urllib.request.urlopen(
|
| 343 |
body = r.read()
|
| 344 |
hv = struct.unpack("<IIII", body[:16])
|
| 345 |
R["anonymous_shard_read"] = {"bytes": len(body), "header": list(hv),
|
|
@@ -348,19 +409,36 @@ def main():
|
|
| 348 |
== man["shards"][0]["sha256"]}
|
| 349 |
except Exception as e:
|
| 350 |
R["anonymous_shard_read"] = {"error": f"{type(e).__name__}: {str(e)[:160]}"}
|
| 351 |
-
R["published_verified"] = bool(not
|
| 352 |
-
and R["anonymous_shard_read"].get("sha256_matches_manifest")
|
|
|
|
| 353 |
print("PUBLISH_JSON_BEGIN")
|
| 354 |
print(json.dumps(R, indent=1, default=str))
|
| 355 |
print("PUBLISH_JSON_END")
|
| 356 |
-
|
| 357 |
-
|
| 358 |
-
|
| 359 |
-
|
| 360 |
-
|
| 361 |
-
|
| 362 |
-
"publish_mix resumes from the ones that landed", flush=True)
|
| 363 |
sys.exit(11)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 364 |
|
| 365 |
|
| 366 |
if __name__ == "__main__":
|
|
|
|
| 120 |
def remote_digests(api, repo):
|
| 121 |
"""path -> content sha256 for what is already on the Hub, so a re-run resumes instead of resending.
|
| 122 |
|
| 123 |
+
Only LFS-tracked files expose a content sha256 (`lfs.oid` is the LFS OID, so sha256 of the payload by
|
| 124 |
+
construction). Small non-LFS files expose only a git blob sha, which is not a content sha256, so they
|
| 125 |
+
are reported None (unknown) and verified by hashing the bytes back from `resolve` instead. `lfs` arrives
|
| 126 |
+
as a dict from huggingface_hub, so attribute access alone silently returns None for every shard and
|
| 127 |
+
turns a resume into a full 2.3 GB re-upload (Gate 2f).
|
| 128 |
+
|
| 129 |
+
A digest that is not 64 hex characters is reported unknown rather than as a mismatch: the field has
|
| 130 |
+
appeared both bare and `sha256:`-prefixed across Hub versions, and reading a format change as a content
|
| 131 |
+
disagreement would re-send everything *and* fail a repo that is actually correct.
|
| 132 |
"""
|
| 133 |
out = {}
|
| 134 |
for ent in api.dataset_info(repo, repo_type="dataset", files_metadata=True).siblings or []:
|
| 135 |
lfs = getattr(ent, "lfs", None)
|
| 136 |
if isinstance(lfs, dict):
|
| 137 |
+
raw = lfs.get("sha256") or lfs.get("oid")
|
| 138 |
elif lfs is not None:
|
| 139 |
+
raw = getattr(lfs, "sha256", None) or getattr(lfs, "oid", None)
|
| 140 |
else:
|
| 141 |
+
raw = None
|
| 142 |
+
if isinstance(raw, str):
|
| 143 |
+
raw = raw.split(":")[-1].strip().lower()
|
| 144 |
+
ok = isinstance(raw, str) and len(raw) == 64 and all(c in "0123456789abcdef" for c in raw)
|
| 145 |
+
out[ent.rfilename] = raw if ok else None
|
| 146 |
return out
|
| 147 |
|
| 148 |
|
|
|
|
| 195 |
f"{str(aud.get('mix_bytes_sha256'))[:16]}, this manifest is {want_digest[:16]}")
|
| 196 |
if int(aud.get("val_documents_audited") or 0) <= 0:
|
| 197 |
problems.append("the audit never scanned the held-out validation shards")
|
| 198 |
+
# A manifest whose records do not add up is not a mix, and every consumer downstream trusts it: the
|
| 199 |
+
# reader builds its corpus from man["shards"], the audit digests it, and the report quotes its totals.
|
| 200 |
+
# Nothing else compares the record list against the claimed counts, so a rebuild that truncated the
|
| 201 |
+
# list while leaving `total_tokens` alone would certify 1.11 B tokens of a 1.03 B corpus (review).
|
| 202 |
+
rec_tokens = sum(int(s.get("tokens") or 0) for s in man.get("shards") or [])
|
| 203 |
+
tot_claim = man.get("total_tokens")
|
| 204 |
+
if len(man.get("shards") or []) != int(man.get("n_shards") or -1):
|
| 205 |
+
problems.append(f"manifest carries {len(man.get('shards') or [])} shard records but claims "
|
| 206 |
+
f"n_shards={man.get('n_shards')}")
|
| 207 |
+
if rec_tokens != int(tot_claim or -1):
|
| 208 |
+
problems.append(f"the shard records sum to {rec_tokens} tokens but total_tokens says "
|
| 209 |
+
f"{tot_claim}")
|
| 210 |
if problems:
|
| 211 |
print("FATAL: contamination gate refuses publication:", flush=True)
|
| 212 |
for p in problems:
|
|
|
|
| 282 |
print(f" (no resumable listing: {type(e).__name__}; uploading everything)", flush=True)
|
| 283 |
|
| 284 |
todo = []
|
| 285 |
+
digests = {}
|
| 286 |
for np, lp in to_upload:
|
| 287 |
h = hashlib.sha256()
|
| 288 |
with open(lp, "rb") as fh:
|
| 289 |
for chunk in iter(lambda: fh.read(1 << 20), b""):
|
| 290 |
h.update(chunk)
|
| 291 |
+
local_sha = digests[np] = h.hexdigest()
|
| 292 |
if remote_sha.get(np) == local_sha:
|
| 293 |
print(f" = {np} already on the Hub with this sha256, skipped", flush=True)
|
| 294 |
continue
|
| 295 |
todo.append((np, lp, local_sha))
|
| 296 |
print(f" {len(todo)} of {len(to_upload)} files still need uploading", flush=True)
|
| 297 |
|
| 298 |
+
# Commit in batches of at most 120 files, not one commit per file and not one commit for everything.
|
| 299 |
+
# Gate 2f proved the first is refused: 139 sequential `upload_file` commits succeeded and then the Hub
|
| 300 |
+
# began rejecting every request -- including 30-byte files -- with a bare HfHubHTTPError, a commit-rate
|
| 301 |
+
# limit rather than a network fault. The second is limited too: `upload_folder` above a ~150-file
|
| 302 |
+
# threshold wants to split into multiple commits, and a cold publish of all 173 files would hit that.
|
| 303 |
failures = []
|
| 304 |
+
CHUNK = 120
|
| 305 |
+
for i in range(0, len(todo), CHUNK):
|
| 306 |
+
batch = todo[i:i + CHUNK]
|
| 307 |
+
paths = [t[0] for t in batch]
|
| 308 |
+
mb = sum(os.path.getsize(p) for _, p, _ in batch) / 1e6
|
| 309 |
try:
|
| 310 |
+
ci = api.upload_folder(repo_id=a.repo, repo_type="dataset", folder_path=a.root,
|
| 311 |
+
path_in_repo="", allow_patterns=paths, max_workers=2,
|
| 312 |
+
commit_message=f"publish mix-v1: batch {i // CHUNK + 1} "
|
| 313 |
+
f"of {len(batch)} file(s)",
|
| 314 |
+
commit_description="batched commits; see audit.json for the "
|
| 315 |
+
"contamination verdict bound to these bytes")
|
| 316 |
+
print(f" + batch {i // CHUNK + 1}: {len(batch)} file(s), {mb:.0f} MB, "
|
| 317 |
+
f"commit {getattr(ci, 'oid', None) or ci}", flush=True)
|
| 318 |
except Exception as e:
|
| 319 |
print(f" batched commit failed ({type(e).__name__}: {str(e)[:160]}); "
|
| 320 |
"falling back to per-file uploads with retries", flush=True)
|
| 321 |
+
for np, lp, local_sha in batch:
|
| 322 |
for attempt in (1, 2, 3):
|
| 323 |
try:
|
| 324 |
with open(lp, "rb") as fh:
|
|
|
|
| 341 |
sys.exit(11)
|
| 342 |
|
| 343 |
# Independent verification from the Hub, not from the client that just wrote it. A batched commit can
|
| 344 |
+
# return without raising and still not carry every file, so presence is checked per path; and because
|
| 345 |
+
# the reader that consumes this repo downloads it anonymously, every path is content-checked -- by the
|
| 346 |
+
# Hub's own digest where it publishes one (LFS), by hashing the bytes back over `resolve` where it does
|
| 347 |
+
# not (the json/md/u8 files, all of which are small). A failed listing is never allowed to read as "no
|
| 348 |
+
# mismatches", which is what `except Exception: after = {}` used to do (review finding).
|
| 349 |
+
after = {}
|
| 350 |
+
for attempt in (1, 2, 3):
|
| 351 |
+
try:
|
| 352 |
+
after = remote_digests(api, a.repo)
|
| 353 |
+
break
|
| 354 |
+
except Exception as e:
|
| 355 |
+
print(f" verification listing attempt {attempt}: {type(e).__name__}: {str(e)[:140]}",
|
| 356 |
+
flush=True)
|
| 357 |
+
__import__("time").sleep(20 * attempt)
|
| 358 |
+
remote = set()
|
| 359 |
+
for attempt in (1, 2, 3):
|
| 360 |
+
try:
|
| 361 |
+
remote = set(api.list_repo_files(repo_id=a.repo, repo_type="dataset"))
|
| 362 |
+
break
|
| 363 |
+
except Exception as e:
|
| 364 |
+
print(f" file listing attempt {attempt}: {type(e).__name__}: {str(e)[:140]}", flush=True)
|
| 365 |
+
__import__("time").sleep(20 * attempt)
|
| 366 |
+
if not remote:
|
| 367 |
+
print("FATAL: the repo cannot be listed after the upload -- refusing to call it published",
|
| 368 |
+
flush=True)
|
| 369 |
+
sys.exit(15)
|
| 370 |
+
|
| 371 |
want = {np for np, _ in to_upload} | {".gitattributes"}
|
| 372 |
local_by_path = {np: lp for np, lp in to_upload}
|
| 373 |
+
base_url = f"https://huggingface.co/datasets/{a.repo}/resolve/main/"
|
| 374 |
+
sha_bad, unresolved = [], []
|
| 375 |
+
for np, lp in to_upload:
|
| 376 |
+
rsha = after.get(np)
|
| 377 |
+
if rsha is None:
|
| 378 |
+
if os.path.getsize(lp) > 8 << 20:
|
| 379 |
+
unresolved.append(f"{np} (too large to read back, no Hub digest)")
|
| 380 |
+
continue
|
| 381 |
+
try:
|
| 382 |
+
body = urllib.request.urlopen(base_url + np, timeout=300).read()
|
| 383 |
+
except Exception as e:
|
| 384 |
+
unresolved.append(f"{np} ({type(e).__name__})")
|
| 385 |
+
continue
|
| 386 |
+
if hashlib.sha256(body).hexdigest() != digests[np]:
|
| 387 |
sha_bad.append(np)
|
| 388 |
+
elif rsha != digests[np]:
|
| 389 |
+
sha_bad.append(np)
|
| 390 |
+
missing = sorted(want - remote)
|
| 391 |
+
# A shard on the Hub that the manifest does not name is not "extra disk", it is a shard some reader
|
| 392 |
+
# globbing `mix/*.bin` will train on while the audit never saw it.
|
| 393 |
+
extra_bin = sorted(p for p in (remote - want) if p.endswith(".bin"))
|
| 394 |
R = {"repo": a.repo, "files_expected": len(want), "files_present": len(want & remote),
|
| 395 |
+
"missing": missing, "sha_mismatch_on_hub": sorted(sha_bad)[:20],
|
| 396 |
+
"extra_shards_on_hub": extra_bin[:20], "unverified": unresolved[:20],
|
| 397 |
+
"sha_verified_via_hub": sum(1 for k, v in after.items() if v and k in local_by_path)}
|
| 398 |
+
print("after upload: %d/%d present, %d missing, %d sha mismatch, %d extra shards, %d unverified" % (
|
| 399 |
+
R["files_present"], R["files_expected"], len(missing), len(sha_bad), len(extra_bin),
|
| 400 |
+
len(unresolved)), flush=True)
|
| 401 |
first = man["shards"][0]["file"]
|
|
|
|
| 402 |
try:
|
| 403 |
+
with urllib.request.urlopen(base_url + first, timeout=300) as r:
|
| 404 |
body = r.read()
|
| 405 |
hv = struct.unpack("<IIII", body[:16])
|
| 406 |
R["anonymous_shard_read"] = {"bytes": len(body), "header": list(hv),
|
|
|
|
| 409 |
== man["shards"][0]["sha256"]}
|
| 410 |
except Exception as e:
|
| 411 |
R["anonymous_shard_read"] = {"error": f"{type(e).__name__}: {str(e)[:160]}"}
|
| 412 |
+
R["published_verified"] = bool(not missing and not sha_bad and not extra_bin and not unresolved
|
| 413 |
+
and R["anonymous_shard_read"].get("sha256_matches_manifest")
|
| 414 |
+
and R["anonymous_shard_read"].get("n_tokens_matches_manifest"))
|
| 415 |
print("PUBLISH_JSON_BEGIN")
|
| 416 |
print(json.dumps(R, indent=1, default=str))
|
| 417 |
print("PUBLISH_JSON_END")
|
| 418 |
+
# Each of these is a different thing to have gone wrong, so each gets its own exit code: 11 a path
|
| 419 |
+
# never landed (re-running resumes), 12 the Hub holds bytes that are not ours, 14 a shard the manifest
|
| 420 |
+
# does not know about, 13 something could not be confirmed at all.
|
| 421 |
+
if missing:
|
| 422 |
+
print(f"FATAL: {len(missing)} file(s) still absent after the commit -- re-running publish_mix "
|
| 423 |
+
"resumes from the ones that landed", flush=True)
|
|
|
|
| 424 |
sys.exit(11)
|
| 425 |
+
if sha_bad:
|
| 426 |
+
print(f"FATAL: {len(sha_bad)} file(s) on the Hub do not match the local bytes: "
|
| 427 |
+
f"{sha_bad[:6]}", flush=True)
|
| 428 |
+
sys.exit(12)
|
| 429 |
+
if extra_bin:
|
| 430 |
+
print(f"FATAL: {len(extra_bin)} .bin file(s) on the Hub are not in this manifest: "
|
| 431 |
+
f"{extra_bin[:6]} -- a reader that globs the directory would train on bytes the audit never "
|
| 432 |
+
"saw, so they must be deleted before this mix is called published", flush=True)
|
| 433 |
+
sys.exit(14)
|
| 434 |
+
if unresolved:
|
| 435 |
+
print(f"FATAL: {len(unresolved)} file(s) could not be content-verified after upload: "
|
| 436 |
+
f"{unresolved[:6]}", flush=True)
|
| 437 |
+
sys.exit(13)
|
| 438 |
+
if not R["published_verified"]:
|
| 439 |
+
print("FATAL: the anonymous read-back of the first shard did not match the manifest -- the repo "
|
| 440 |
+
"is not provably serving the bytes it claims", flush=True)
|
| 441 |
+
sys.exit(13)
|
| 442 |
|
| 443 |
|
| 444 |
if __name__ == "__main__":
|