build_mix: hub-init failure is fatal (not silent), empty merge exits 4; add token-scope probe
Browse files- build/build_mix.py +20 -9
- probes/p2_token_scope.py +130 -0
build/build_mix.py
CHANGED
|
@@ -219,7 +219,8 @@ class State:
|
|
| 219 |
|
| 220 |
# --------------------------------------------------------------------- stage
|
| 221 |
def stage(args, st, tokenizer):
|
| 222 |
-
remote = HS.remote_stage_files(args.hub_repo) if (
|
|
|
|
| 223 |
if args.hub_repo:
|
| 224 |
print(f"hub: {len(remote)} staged files already on {args.hub_repo}", flush=True)
|
| 225 |
for src in SOURCES:
|
|
@@ -228,7 +229,7 @@ def stage(args, st, tokenizer):
|
|
| 228 |
if rec.get("complete"):
|
| 229 |
print(f"stage: skip {key} (complete)", flush=True)
|
| 230 |
continue
|
| 231 |
-
if args.hub_repo:
|
| 232 |
# A source whose shards AND record are already on the Hub is done; pulling the record is
|
| 233 |
# cheaper than re-staging and keeps the merge's proportional weights honest.
|
| 234 |
have = [p for p in remote if p.startswith(f"stage/{key}/") and p.endswith(".bin")]
|
|
@@ -352,7 +353,7 @@ def stage(args, st, tokenizer):
|
|
| 352 |
st.commit(new_h, new_k)
|
| 353 |
# /kaggle/working is destroyed when the session ends, so a source that took 40 minutes to stage
|
| 354 |
# must not exist only there. Published per source, because that is the unit a resume can skip.
|
| 355 |
-
if complete and args.hub_repo:
|
| 356 |
with open(os.path.join(st.stage_dir, key, "record.json"), "w") as f:
|
| 357 |
json.dump({"key": key, "src": {k: v for k, v in src.items()},
|
| 358 |
"shards": shards, "tokens": got, "docs_seen": docs_seen,
|
|
@@ -392,8 +393,11 @@ def merge(args, st, tokenizer):
|
|
| 392 |
streams.append({"key": key, "weight": float(weight), "credit": 0.0, "it": gen(iters),
|
| 393 |
"pending": None, "emitted_tokens": 0, "exhausted": False})
|
| 394 |
if not streams:
|
| 395 |
-
|
| 396 |
-
|
|
|
|
|
|
|
|
|
|
| 397 |
total_w = sum(s["weight"] for s in streams)
|
| 398 |
print(f"merge: {len(streams)} sources, "
|
| 399 |
f"{sum(s['weight'] for s in streams) / 1e6:,.0f}M staged tokens", flush=True)
|
|
@@ -491,15 +495,21 @@ def main():
|
|
| 491 |
|
| 492 |
global HS
|
| 493 |
HS = None
|
|
|
|
|
|
|
|
|
|
| 494 |
if args.hub_repo:
|
| 495 |
try:
|
| 496 |
import hubsync as HS_mod
|
| 497 |
HS = HS_mod
|
| 498 |
HS.ensure_repo(args.hub_repo)
|
| 499 |
except Exception as e:
|
| 500 |
-
print(f"
|
| 501 |
-
|
| 502 |
-
|
|
|
|
|
|
|
|
|
|
| 503 |
|
| 504 |
st = State(args.root)
|
| 505 |
if args.hub_repo and not st.d["staged"]:
|
|
@@ -522,7 +532,8 @@ def main():
|
|
| 522 |
print(f"hub: restored {len(st.d['staged'])} staged source(s) "
|
| 523 |
f"({st.staged_tokens_total():,} tokens)", flush=True)
|
| 524 |
except Exception as e:
|
| 525 |
-
print(f"
|
|
|
|
| 526 |
|
| 527 |
t0 = time.monotonic()
|
| 528 |
if not args.merge_only:
|
|
|
|
| 219 |
|
| 220 |
# --------------------------------------------------------------------- stage
|
| 221 |
def stage(args, st, tokenizer):
|
| 222 |
+
remote = HS.remote_stage_files(args.hub_repo) if (HS and args.hub_repo
|
| 223 |
+
and not args.merge_only) else set()
|
| 224 |
if args.hub_repo:
|
| 225 |
print(f"hub: {len(remote)} staged files already on {args.hub_repo}", flush=True)
|
| 226 |
for src in SOURCES:
|
|
|
|
| 229 |
if rec.get("complete"):
|
| 230 |
print(f"stage: skip {key} (complete)", flush=True)
|
| 231 |
continue
|
| 232 |
+
if HS and args.hub_repo:
|
| 233 |
# A source whose shards AND record are already on the Hub is done; pulling the record is
|
| 234 |
# cheaper than re-staging and keeps the merge's proportional weights honest.
|
| 235 |
have = [p for p in remote if p.startswith(f"stage/{key}/") and p.endswith(".bin")]
|
|
|
|
| 353 |
st.commit(new_h, new_k)
|
| 354 |
# /kaggle/working is destroyed when the session ends, so a source that took 40 minutes to stage
|
| 355 |
# must not exist only there. Published per source, because that is the unit a resume can skip.
|
| 356 |
+
if complete and HS and args.hub_repo:
|
| 357 |
with open(os.path.join(st.stage_dir, key, "record.json"), "w") as f:
|
| 358 |
json.dump({"key": key, "src": {k: v for k, v in src.items()},
|
| 359 |
"shards": shards, "tokens": got, "docs_seen": docs_seen,
|
|
|
|
| 393 |
streams.append({"key": key, "weight": float(weight), "credit": 0.0, "it": gen(iters),
|
| 394 |
"pending": None, "emitted_tokens": 0, "exhausted": False})
|
| 395 |
if not streams:
|
| 396 |
+
# The cold-resume rehearsal exited 0 here while producing an empty mix, and verify_mix then
|
| 397 |
+
# reported recount_matches_manifest: true because 0 == 0. An empty mix is the worst possible
|
| 398 |
+
# silent success in this pipeline, so it is a failure and not a return.
|
| 399 |
+
print("FATAL merge: nothing staged. Refusing to emit an empty mix.", flush=True)
|
| 400 |
+
sys.exit(4)
|
| 401 |
total_w = sum(s["weight"] for s in streams)
|
| 402 |
print(f"merge: {len(streams)} sources, "
|
| 403 |
f"{sum(s['weight'] for s in streams) / 1e6:,.0f}M staged tokens", flush=True)
|
|
|
|
| 495 |
|
| 496 |
global HS
|
| 497 |
HS = None
|
| 498 |
+
# Sync is best-effort to SET UP but not to USE: if the Hub cannot be reached or written, turning it
|
| 499 |
+
# off silently would reintroduce the exact failure this exists to prevent (a killed session losing
|
| 500 |
+
# hours of streaming). So a configured repo that cannot be initialised is fatal unless --no-hub.
|
| 501 |
if args.hub_repo:
|
| 502 |
try:
|
| 503 |
import hubsync as HS_mod
|
| 504 |
HS = HS_mod
|
| 505 |
HS.ensure_repo(args.hub_repo)
|
| 506 |
except Exception as e:
|
| 507 |
+
print(f"FATAL hub sync could not initialise: {type(e).__name__}: {str(e)[:200]}",
|
| 508 |
+
flush=True)
|
| 509 |
+
print("Refusing to run a multi-hour build that is not interruption-safe. Resolve the Hub "
|
| 510 |
+
"write path, or pass --hub-repo '' only if you accept losing all staged work when the "
|
| 511 |
+
"session ends.", flush=True)
|
| 512 |
+
sys.exit(3)
|
| 513 |
|
| 514 |
st = State(args.root)
|
| 515 |
if args.hub_repo and not st.d["staged"]:
|
|
|
|
| 532 |
print(f"hub: restored {len(st.d['staged'])} staged source(s) "
|
| 533 |
f"({st.staged_tokens_total():,} tokens)", flush=True)
|
| 534 |
except Exception as e:
|
| 535 |
+
print(f"FATAL hub restore failed: {type(e).__name__}: {str(e)[:200]}", flush=True)
|
| 536 |
+
sys.exit(3)
|
| 537 |
|
| 538 |
t0 = time.monotonic()
|
| 539 |
if not args.merge_only:
|
probes/p2_token_scope.py
ADDED
|
@@ -0,0 +1,130 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
# Phase 2 blocker check: can this token write DATASETS, or only models?
|
| 2 |
+
#
|
| 3 |
+
# Why this exists. `p1-credential-roundtrip` proved a job can create and delete a **model** repo. The
|
| 4 |
+
# cold-resume rehearsal then failed to create a **dataset** repo with `401 Unauthorized` from
|
| 5 |
+
# POST /api/repos/create, using the same token and the same code path. Phase 2 requires publishing the mix
|
| 6 |
+
# as a public `ounce100m-*` *dataset*, and the stage repo is a dataset too -- so if the hypothesis is
|
| 7 |
+
# right, this is a hard blocker on Gate 2 that only the user can clear by adjusting the token.
|
| 8 |
+
#
|
| 9 |
+
# The token is fine-grained (`repo.write` scoped to the user). HF splits "write to models" and "write to
|
| 10 |
+
# datasets" into separate scopes, and `repo.write` has historically meant the former, so the hypothesis is
|
| 11 |
+
# plausible rather than a mystery.
|
| 12 |
+
#
|
| 13 |
+
# Test design: for each repo type, create → upload → read back anonymously → delete, and report the
|
| 14 |
+
# outcome per step. Anything created here is deleted in the same run; nothing lingers in the account.
|
| 15 |
+
# It also dumps the token's own declared capabilities, which is the authoritative answer rather than an
|
| 16 |
+
# inference from a 401.
|
| 17 |
+
#
|
| 18 |
+
# CPU session, zero GPU quota. No credentials in this file (D-006 helper).
|
| 19 |
+
|
| 20 |
+
import json
|
| 21 |
+
import os
|
| 22 |
+
import urllib.error
|
| 23 |
+
import urllib.request
|
| 24 |
+
|
| 25 |
+
import ounce100m_credentials
|
| 26 |
+
|
| 27 |
+
R = {}
|
| 28 |
+
STAMP = "scopetest"
|
| 29 |
+
NAMES = {"model": f"Cion-lab/ounce100m-{STAMP}-model-DELETEME",
|
| 30 |
+
"dataset": f"Cion-lab/ounce100m-{STAMP}-dataset-DELETEME"}
|
| 31 |
+
|
| 32 |
+
|
| 33 |
+
def guard(name, fn):
|
| 34 |
+
try:
|
| 35 |
+
R[name] = fn()
|
| 36 |
+
except Exception as e:
|
| 37 |
+
detail = ""
|
| 38 |
+
resp = getattr(e, "response", None)
|
| 39 |
+
if resp is not None:
|
| 40 |
+
try:
|
| 41 |
+
detail = resp.text[:400]
|
| 42 |
+
except Exception:
|
| 43 |
+
pass
|
| 44 |
+
R[name] = {"error": f"{type(e).__name__}: {str(e)[:280]}", "body": detail}
|
| 45 |
+
|
| 46 |
+
|
| 47 |
+
guard("install", lambda: ounce100m_credentials.install(verify=True))
|
| 48 |
+
TOKEN = os.environ.get("HF_TOKEN", "")
|
| 49 |
+
|
| 50 |
+
|
| 51 |
+
def whoami_full():
|
| 52 |
+
req = urllib.request.Request("https://huggingface.co/whoami-v2",
|
| 53 |
+
headers={"Authorization": f"Bearer {TOKEN}"})
|
| 54 |
+
with urllib.request.urlopen(req, timeout=45) as r:
|
| 55 |
+
d = json.loads(r.read().decode())
|
| 56 |
+
# report capability *shape*, never the token itself
|
| 57 |
+
out = {"name": d.get("name"), "type": d.get("type")}
|
| 58 |
+
full = d.get("full") or {}
|
| 59 |
+
out["full_keys"] = sorted(full.keys())
|
| 60 |
+
auth = full.get("auth") or d.get("auth") or {}
|
| 61 |
+
out["auth"] = {k: v for k, v in auth.items() if k != "accessToken"}
|
| 62 |
+
for k in ("capabilities", "scopes", "permissions"):
|
| 63 |
+
if k in full:
|
| 64 |
+
out[k] = full[k]
|
| 65 |
+
elif k in d:
|
| 66 |
+
out[k] = d[k]
|
| 67 |
+
if "accessToken" in full:
|
| 68 |
+
at = full["accessToken"]
|
| 69 |
+
out["accessToken_type"] = type(at).__name__
|
| 70 |
+
if isinstance(at, dict):
|
| 71 |
+
out["accessToken_scopes"] = at.get("scopes") or at.get("permission")
|
| 72 |
+
return out
|
| 73 |
+
|
| 74 |
+
|
| 75 |
+
guard("whoami_v2", whoami_full)
|
| 76 |
+
|
| 77 |
+
|
| 78 |
+
def trial(kind):
|
| 79 |
+
from huggingface_hub import HfApi, create_repo, delete_repo
|
| 80 |
+
|
| 81 |
+
repo = NAMES[kind]
|
| 82 |
+
api = HfApi(token=TOKEN)
|
| 83 |
+
steps = {}
|
| 84 |
+
create_repo(repo_id=repo, repo_type=kind, private=False, exist_ok=True)
|
| 85 |
+
steps["create"] = "ok"
|
| 86 |
+
payload = f"ounce100m scope probe {kind}\n".encode()
|
| 87 |
+
api.upload_file(path_or_fileobj=payload, path_in_repo="probe.txt", repo_id=repo,
|
| 88 |
+
repo_type=kind, commit_message="scope probe")
|
| 89 |
+
steps["upload"] = "ok"
|
| 90 |
+
url = (f"https://huggingface.co/{repo}/resolve/main/probe.txt" if kind == "model"
|
| 91 |
+
else f"https://huggingface.co/datasets/{repo}/resolve/main/probe.txt")
|
| 92 |
+
try:
|
| 93 |
+
with urllib.request.urlopen(url, timeout=60) as r:
|
| 94 |
+
steps["anonymous_readback_matches"] = r.read().strip() == payload.strip()
|
| 95 |
+
except urllib.error.HTTPError as e:
|
| 96 |
+
steps["anonymous_readback"] = f"HTTP {e.code}"
|
| 97 |
+
delete_repo(repo_id=repo, repo_type=kind)
|
| 98 |
+
steps["delete"] = "ok"
|
| 99 |
+
return steps
|
| 100 |
+
|
| 101 |
+
|
| 102 |
+
for kind in ("model", "dataset"):
|
| 103 |
+
guard(f"trial_{kind}", lambda k=kind: trial(k))
|
| 104 |
+
|
| 105 |
+
# confirm the deletes really took, so nothing is left behind in the account
|
| 106 |
+
def listing():
|
| 107 |
+
out = {}
|
| 108 |
+
for kind, q in (("model", "models"), ("dataset", "datasets")):
|
| 109 |
+
try:
|
| 110 |
+
with urllib.request.urlopen(
|
| 111 |
+
f"https://huggingface.co/api/{q}?author=Cion-lab", timeout=45) as r:
|
| 112 |
+
out[kind] = sorted(m["id"] for m in json.loads(r.read().decode()))
|
| 113 |
+
except Exception as e:
|
| 114 |
+
out[kind] = f"{type(e).__name__}"
|
| 115 |
+
return out
|
| 116 |
+
|
| 117 |
+
|
| 118 |
+
guard("namespace_after", listing)
|
| 119 |
+
|
| 120 |
+
R["verdict"] = {
|
| 121 |
+
"model_write": "ok" if isinstance(R.get("trial_model"), dict) and
|
| 122 |
+
R["trial_model"].get("create") == "ok" else R.get("trial_model"),
|
| 123 |
+
"dataset_write": "ok" if isinstance(R.get("trial_dataset"), dict) and
|
| 124 |
+
R["trial_dataset"].get("create") == "ok" else R.get("trial_dataset"),
|
| 125 |
+
}
|
| 126 |
+
R["BLOCKS_GATE_2"] = R["verdict"]["dataset_write"] != "ok"
|
| 127 |
+
|
| 128 |
+
print("SCOPE_JSON_BEGIN")
|
| 129 |
+
print(json.dumps(R, indent=1, default=str)[:5000])
|
| 130 |
+
print("SCOPE_JSON_END")
|