Buckets:
| """Convert the ClawHub public-API crawl (code/clawhub_api.py output) into the processed layout and ship it: | |
| processed/clawhub/api_listing.parquet every public skill: stats (downloads/installs/stars), versions, topics | |
| processed/clawhub/api_delta_skills.parquet skills created / re-versioned after the HF dump, from their zip bundles | |
| processed/clawhub/api_delta_files.parquet the other files of those bundles | |
| raw/clawhub_api/{listing.jsonl, zips/} raw crawl | |
| Zips contain the published bundle plus `_meta.json` and a registry-generated `skill-card.md` (both kept as files, | |
| flagged generated=True). Waits for the crawler service to finish first. | |
| """ | |
| import json, sys, time, zipfile, subprocess, os, gc, ctypes | |
| from pathlib import Path | |
| import pyarrow as pa, pyarrow.parquet as pq | |
| sys.path.insert(0, str(Path(__file__).parent)) | |
| from lib_skill import parse_skill_md, cjk_ratio | |
| ROOT = Path(__file__).resolve().parent.parent | |
| RAW = ROOT / "data/raw/clawhub_api" | |
| OUT = ROOT / "work/clawhub_api"; OUT.mkdir(parents=True, exist_ok=True) | |
| B = "hf://buckets/Mercity/SkillsStorage" | |
| HF = str(ROOT / ".venv/bin/hf") | |
| GENERATED = {"_meta.json", "skill-card.md"} | |
| DONE = ROOT / "logs/build_clawhub_api.done" | |
| if DONE.exists(): print("already done", flush=True); sys.exit(0) | |
| while subprocess.run(["systemctl", "is-active", "--quiet", "skills-clawhub-api"]).returncode == 0: | |
| time.sleep(60) | |
| if not (RAW / "listing.jsonl").exists(): sys.exit("listing not finished") | |
| def cp(src, dst): | |
| for i in range(6): | |
| if subprocess.run([HF, "buckets", "cp", str(src), dst], capture_output=True).returncode == 0: return | |
| time.sleep(15 * (i + 1)) | |
| raise RuntimeError(f"upload failed: {src}") | |
| FINAL = [OUT / "api_listing.parquet", OUT / "api_delta_skills.parquet", OUT / "api_delta_files.parquet"] | |
| def valid(p): | |
| try: pq.ParquetFile(p); return True | |
| except Exception: return False | |
| if not all(p.exists() and valid(p) for p in FINAL): # restart-safe: finished outputs are never rebuilt | |
| # listing -> stats table (dedup by owner/slug, keep latest seen) | |
| items = {} | |
| for line in open(RAW / "listing.jsonl"): | |
| it = json.loads(line); items[(it["ownerHandle"], it["slug"])] = it | |
| rows = [] | |
| for (owner, slug), it in items.items(): | |
| st, lv = it.get("stats") or {}, it.get("latestVersion") or {} | |
| rows.append({"owner": owner, "slug": slug, "display_name": it.get("displayName"), "summary": it.get("summary"), | |
| "topics": it.get("topics") or [], "downloads": st.get("downloads"), "installs": st.get("installs"), | |
| "stars": st.get("stars"), "comments": st.get("comments"), "n_versions": st.get("versions"), | |
| "latest_version": lv.get("version"), "latest_version_at": lv.get("createdAt"), "license": lv.get("license"), | |
| "created_at": it.get("createdAt"), "updated_at": it.get("updatedAt"), | |
| "metadata": json.dumps(it.get("metadata"), ensure_ascii=False) if it.get("metadata") else None}) | |
| pq.write_table(pa.Table.from_pylist(rows), str(FINAL[0]) + ".tmp", compression="zstd") | |
| print("listing rows", len(rows), flush=True) | |
| by_key = {f"{o}__{s}": (it.get("latestVersion") or {}).get("version") for (o, s), it in items.items()} | |
| del items, rows; gc.collect() | |
| # zips -> delta skills + files (files streamed to disk in batches) | |
| sk, fi, fi_w, n_fi = [], [], None, 0 | |
| FI_SCHEMA = pa.schema([("skill_uid", pa.string()), ("source", pa.string()), ("rel_path", pa.string()), ("size", pa.int64()), | |
| ("content", pa.string()), ("content_bin", pa.binary()), ("has_content", pa.bool_()), ("generated", pa.bool_())]) | |
| for z in sorted((RAW / "zips").glob("*.zip")): | |
| zf = zipfile.ZipFile(z) | |
| meta = json.loads(zf.read("_meta.json")) if "_meta.json" in zf.namelist() else {} | |
| owner, slug = z.stem.split("__", 1) | |
| uid = f"clawhub:{owner}/{slug}@{meta.get('version') or by_key.get(z.stem)}" | |
| md = zf.read("SKILL.md").decode("utf-8", "replace") if "SKILL.md" in zf.namelist() else "" | |
| ok, name, desc, _, body = parse_skill_md(md) | |
| others = [i for i in zf.infolist() if not i.is_dir() and i.filename != "SKILL.md"] | |
| sk.append({"skill_uid": uid, "source": "clawhub_api", "owner": owner, "slug": slug, "version": meta.get("version"), | |
| "published_at": meta.get("publishedAt"), "skill_md": md, "name": name, "description": desc, | |
| "frontmatter_valid": ok, "body_chars": len(body), "cjk_ratio": cjk_ratio(md), "license": "MIT-0", | |
| "n_bundle_files": sum(1 for i in others if i.filename not in GENERATED), | |
| "bundle_bytes": sum(i.file_size for i in others if i.filename not in GENERATED)}) | |
| for i in others: | |
| data = zf.read(i); text = b"\0" not in data[:8192] | |
| fi.append({"skill_uid": uid, "source": "clawhub_api", "rel_path": i.filename, "size": i.file_size, | |
| "content": data.decode("utf-8", "replace") if text else None, "content_bin": None if text else data, | |
| "has_content": True, "generated": i.filename in GENERATED}) | |
| if len(fi) >= 2000: | |
| fi_w = fi_w or pq.ParquetWriter(str(FINAL[2]) + ".tmp", FI_SCHEMA, compression="zstd") | |
| fi_w.write_table(pa.Table.from_pylist(fi, FI_SCHEMA)); n_fi += len(fi); fi = [] | |
| fi_w = fi_w or pq.ParquetWriter(str(FINAL[2]) + ".tmp", FI_SCHEMA, compression="zstd") | |
| if fi: fi_w.write_table(pa.Table.from_pylist(fi, FI_SCHEMA)); n_fi += len(fi) | |
| fi_w.close() | |
| pq.write_table(pa.Table.from_pylist(sk), str(FINAL[1]) + ".tmp", compression="zstd") | |
| print("delta skills", len(sk), "files", n_fi, flush=True) | |
| for p in FINAL: os.replace(str(p) + ".tmp", p) | |
| del sk, fi; gc.collect() | |
| try: ctypes.CDLL("libc.so.6").malloc_trim(0) # hand memory back before the upload tool runs in this cgroup | |
| except OSError: pass | |
| for p in OUT.glob("*.parquet"): cp(p, f"{B}/processed/clawhub/{p.name}") | |
| cp(RAW / "listing.jsonl", f"{B}/raw/clawhub_api/listing.jsonl") | |
| subprocess.run([HF, "buckets", "sync", str(RAW / "zips"), f"{B}/raw/clawhub_api/zips"], check=True, capture_output=True) | |
| print("uploaded", flush=True) | |
| DONE.write_text("ok\n") | |
Xet Storage Details
- Size:
- 6.28 kB
- Xet hash:
- 01dbbf688381a5a66897125ea82f31cf777992c58fb12fa6ff1cd370e7fb1ec8
·
Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.