Buckets:

Mercity/SkillsStorage / code /build_clawhub_api.py
Pranav2748's picture
download
raw
6.28 kB
"""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.