pin: kernels/gate2g_publish.py (RepoMissing resume branch, disk gate, PHASE4_PREP_ONLY)
Browse files- kernels/gate2g_publish.py +143 -0
kernels/gate2g_publish.py
ADDED
|
@@ -0,0 +1,143 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
# Gate 2g: rebuild the filtered mix, re-audit it, verify, and PUBLISH mix-v1 with the batched uploader.
|
| 2 |
+
# Gate 2f rebuilt and audited the mix cleanly and then stopped at HfHubHTTPError after ~139 sequential
|
| 3 |
+
# commits: a Hub commit-RATE ceiling, not a network fault (E-033). publish_mix now sends the whole batch
|
| 4 |
+
# as ONE upload_folder commit, resumes by content sha256 (the listing reads lfs.oid, which is a dict --
|
| 5 |
+
# attribute access silently returned None for every shard and turned the "resume" into a 2.3 GB re-send),
|
| 6 |
+
# and hard-exits if any path is still absent or disagrees in sha after the commit.
|
| 7 |
+
import hashlib, json, os, shutil, subprocess, sys, threading, time
|
| 8 |
+
|
| 9 |
+
os.chdir("/kaggle/working")
|
| 10 |
+
sys.path.insert(0, "/kaggle/working")
|
| 11 |
+
REV = "292c9d28bb81251ce09c78a8e84bdbe69032c28e"
|
| 12 |
+
WANT = {
|
| 13 |
+
"ounce100m_credentials.py": ("ounce100m_credentials.py",
|
| 14 |
+
"6525f62f03f2d73650a1eb4f70fcb52d1194caad4ca88b2d8bd8fd54f88339b6"),
|
| 15 |
+
"hubsync.py": ("build/hubsync.py",
|
| 16 |
+
"e9696dd0c48bfcdb3bf201cc6b065d865dab988510646fbefefb49a4b35332c1"),
|
| 17 |
+
"build_mix.py": ("build/build_mix.py",
|
| 18 |
+
"f0cf56256f7dffa2caac4105d59e261d775eb9c50437f142b49ce8ac18db86a1"),
|
| 19 |
+
"verify_mix.py": ("build/verify_mix.py",
|
| 20 |
+
"38a6bb00a42f80cfdd85081847237edbdab22340a81d33f90e2956485ca55c01"),
|
| 21 |
+
"audit_contamination.py": ("build/audit_contamination.py",
|
| 22 |
+
"58656d68977d858d763322b6d7966adc637457c4975e132df89586f3ac4d3088"),
|
| 23 |
+
"publish_mix.py": ("build/publish_mix.py",
|
| 24 |
+
"2d609bb453ed73d9766d958268c858c6f5afed93d758ced8e041054a4aa6f697"),
|
| 25 |
+
}
|
| 26 |
+
BASE = "https://huggingface.co/Cion-lab/ounce100m-code/resolve/" + REV
|
| 27 |
+
for p, (rp, want) in sorted(WANT.items()):
|
| 28 |
+
r = subprocess.run(["curl", "-sfL", BASE + "/" + rp, "-o", p], capture_output=True, text=True)
|
| 29 |
+
assert r.returncode == 0, ("fetch failed", rp, (r.stderr or "")[:200])
|
| 30 |
+
got = hashlib.sha256(open(p, "rb").read()).hexdigest()
|
| 31 |
+
assert got == want, ("SHA MISMATCH", rp, got, want)
|
| 32 |
+
print("OK", p, got[:12], flush=True)
|
| 33 |
+
|
| 34 |
+
import ounce100m_credentials as C
|
| 35 |
+
print("creds:", json.dumps(C.install(verify=True)), flush=True)
|
| 36 |
+
ROOT = "/kaggle/working/mixroot"
|
| 37 |
+
CACHE = "/kaggle/working/refcache"
|
| 38 |
+
FILT = os.path.join(ROOT, "filter")
|
| 39 |
+
os.makedirs(CACHE, exist_ok=True); os.makedirs(FILT, exist_ok=True)
|
| 40 |
+
STAGE_REPO = "Cion-lab/ounce100m-mix-stage"
|
| 41 |
+
MIX_REPO = "Cion-lab/ounce100m-mix-v1"
|
| 42 |
+
T0 = time.time()
|
| 43 |
+
|
| 44 |
+
|
| 45 |
+
def run(argv, label, timeout):
|
| 46 |
+
"""Stream a stage's output line by line: the Gate 2f version buffered everything, so a session that
|
| 47 |
+
hit the notebook timeout reported nothing about the stage that was running."""
|
| 48 |
+
print("=== " + label, flush=True)
|
| 49 |
+
t0 = time.time()
|
| 50 |
+
e = dict(os.environ); e["PYTHONPATH"] = "/kaggle/working"
|
| 51 |
+
killed = []
|
| 52 |
+
p = subprocess.Popen(argv, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True,
|
| 53 |
+
env=e, bufsize=1)
|
| 54 |
+
timer = threading.Timer(timeout, lambda: (killed.append(True), p.kill()))
|
| 55 |
+
timer.daemon = True
|
| 56 |
+
timer.start()
|
| 57 |
+
tail = []
|
| 58 |
+
try:
|
| 59 |
+
for line in p.stdout:
|
| 60 |
+
line = line.rstrip("\n")
|
| 61 |
+
tail.append(line)
|
| 62 |
+
del tail[:-80]
|
| 63 |
+
print(" |", line[:280], flush=True)
|
| 64 |
+
finally:
|
| 65 |
+
timer.cancel()
|
| 66 |
+
rc = p.wait()
|
| 67 |
+
if killed:
|
| 68 |
+
print(" TIMEOUT after %d s, killed" % timeout, flush=True)
|
| 69 |
+
rc = -9
|
| 70 |
+
print("%s_RC %s seconds %.1f elapsed %.0f" % (label, rc, time.time() - t0, time.time() - T0),
|
| 71 |
+
flush=True)
|
| 72 |
+
return rc, "\n".join(tail)
|
| 73 |
+
|
| 74 |
+
|
| 75 |
+
from huggingface_hub import HfApi, hf_hub_download, list_repo_files
|
| 76 |
+
have = {f for f in list_repo_files(STAGE_REPO, repo_type="dataset") if f.startswith("auditref/")}
|
| 77 |
+
for n in ("ref.npy", "refmask.npy", "refmeta.json", "refcounts.json", "errors.json", "srv.json"):
|
| 78 |
+
if "auditref/" + n in have:
|
| 79 |
+
shutil.copyfile(hf_hub_download(STAGE_REPO, "auditref/" + n, repo_type="dataset",
|
| 80 |
+
cache_dir="/tmp/hf"), os.path.join(CACHE, n))
|
| 81 |
+
want_masks = sorted(f for f in have if f.startswith("auditref/filter/") and f.endswith(".u8"))
|
| 82 |
+
for f in want_masks + (["auditref/filter/filter.json"] if "auditref/filter/filter.json" in have else []):
|
| 83 |
+
shutil.copyfile(hf_hub_download(STAGE_REPO, f, repo_type="dataset", cache_dir="/tmp/hf"),
|
| 84 |
+
os.path.join(FILT, f.split("/")[-1]))
|
| 85 |
+
nm = len([x for x in os.listdir(FILT) if x.endswith(".u8")])
|
| 86 |
+
print("masks restored:", nm, "refcache files:", len(os.listdir(CACHE)), flush=True)
|
| 87 |
+
if nm < 15 or not os.path.exists(os.path.join(FILT, "filter.json")):
|
| 88 |
+
print("VERDICT GATE2G_STOP incomplete mask archive", flush=True)
|
| 89 |
+
raise SystemExit(3)
|
| 90 |
+
|
| 91 |
+
rcm, _ = run([sys.executable, "build_mix.py", "--root", ROOT, "--merge-only", "--hub-repo",
|
| 92 |
+
STAGE_REPO, "--exclude-dir", FILT, "--rebuild-merged"], "MERGE_FILTERED", 3600)
|
| 93 |
+
man = json.load(open(os.path.join(ROOT, "manifest.json")))
|
| 94 |
+
print("MANIFEST", json.dumps({k: man.get(k) for k in
|
| 95 |
+
("total_tokens", "val_tokens", "n_shards", "distinct_sources_per_shard_min",
|
| 96 |
+
"contamination_filter")}, default=str)[:900], flush=True)
|
| 97 |
+
# Gate 2e measured this exact figure after dropping 5,878 contaminated documents. The merge is a
|
| 98 |
+
# deterministic round-robin over the same staged shards with the same masks, so equal here means the
|
| 99 |
+
# rebuild is byte-for-byte the mix that was audited; different means something in the inputs moved.
|
| 100 |
+
print("REMERGE_MATCHES_2E", man.get("total_tokens") == 1109714831, man.get("total_tokens"), flush=True)
|
| 101 |
+
if (man.get("n_shards") or 0) < 100 or (man.get("total_tokens") or 0) < 1_000_000_000:
|
| 102 |
+
print("VERDICT GATE2G_STOP rebuild did not produce a real mix", flush=True)
|
| 103 |
+
raise SystemExit(4)
|
| 104 |
+
|
| 105 |
+
rcb, _ = run([sys.executable, "audit_contamination.py", "--root", ROOT, "--ref-cache", CACHE,
|
| 106 |
+
"--out", os.path.join(ROOT, "audit.json")], "AUDIT", 5400)
|
| 107 |
+
aud = json.load(open(os.path.join(ROOT, "audit.json")))
|
| 108 |
+
print("AUDIT_SUMMARY", json.dumps({k: aud.get(k) for k in
|
| 109 |
+
("overlap_total", "mix_documents_audited", "mix_documents_with_any_hit",
|
| 110 |
+
"mix_documents_majority_hit", "val_documents_audited", "val_hits", "tasks_covered",
|
| 111 |
+
"reference_loader_problems", "mix_bytes_sha256", "MEASURED", "AUDIT_PASSED")},
|
| 112 |
+
default=str)[:900], flush=True)
|
| 113 |
+
rcv, _ = run([sys.executable, "verify_mix.py", "--root", ROOT], "VERIFY", 1800)
|
| 114 |
+
if aud.get("AUDIT_PASSED") is not True or rcv != 0:
|
| 115 |
+
print("VERDICT GATE2G_STOP audit or verify failed -- not publishing", flush=True)
|
| 116 |
+
raise SystemExit(5)
|
| 117 |
+
|
| 118 |
+
rcp, _ = run([sys.executable, "publish_mix.py", "--root", ROOT, "--repo", MIX_REPO], "PUBLISH", 5400)
|
| 119 |
+
print("PUBLISH_RC", rcp, flush=True)
|
| 120 |
+
# Read the published bytes back over the public resolve endpoint with no token in the process
|
| 121 |
+
# environment at all: that is what "the Hub serves this mix to anyone" means, and it is the only check
|
| 122 |
+
# that catches a commit that reported success but left a path absent.
|
| 123 |
+
ANON = """
|
| 124 |
+
import hashlib, json, os, urllib.request
|
| 125 |
+
man = json.load(open("@ROOT@/manifest.json"))
|
| 126 |
+
def read(path, want_sha):
|
| 127 |
+
d = urllib.request.urlopen("@MIX@/resolve/main/" + path, timeout=600).read()
|
| 128 |
+
ok = hashlib.sha256(d).hexdigest() == want_sha
|
| 129 |
+
print("ANON_READBACK", path, len(d), "sha_match", ok, flush=True)
|
| 130 |
+
assert ok, path
|
| 131 |
+
idxs = [0, 7, len(man["shards"]) // 2, len(man["shards"]) - 1]
|
| 132 |
+
for i in idxs:
|
| 133 |
+
read(man["shards"][i]["file"], man["shards"][i]["sha256"])
|
| 134 |
+
for v in (man.get("val_shards") or [])[:2]:
|
| 135 |
+
read(v["file"], v["sha256"])
|
| 136 |
+
import huggingface_hub as hh
|
| 137 |
+
print("REMOTE_FILE_COUNT", len(hh.list_repo_files("@REPO@", repo_type="dataset")), flush=True)
|
| 138 |
+
""".replace("@ROOT@", ROOT).replace("@REPO@", MIX_REPO) \
|
| 139 |
+
.replace("@MIX@", "https://huggingface.co/datasets/" + MIX_REPO)
|
| 140 |
+
open("anonread.py", "w").write(ANON)
|
| 141 |
+
rcz, _ = run([sys.executable, "anonread.py"], "ANON_VERIFY", 1800)
|
| 142 |
+
print("VERDICT GATE2G_DONE publish_rc", rcp, "anon_rc", rcz, "tokens", man.get("total_tokens"),
|
| 143 |
+
"overlap", aud.get("overlap_total"), "seconds", round(time.time() - T0, 1), flush=True)
|