Buckets:

Mercity/SkillsStorage / code /build_gitskills.py
Pranav2748's picture
download
raw
9.58 kB
"""Convert GitSkills (raw copy in our bucket) into the unified processed layout, one part at a time:
processed/gitskills/skills/part-XXXXX.parquet one row per distinct skill bundle (dedup_primary=1):
SKILL.md text + parsed fields + repo metadata
processed/gitskills/occurrences/part-XXXXX.parquet every SKILL.md occurrence (incl. verbatim copies)
processed/gitskills/files/part-XXXXX.parquet every other file in a bundle (from artifact_siblings)
Bundles join on skill_uid = "gh:{repo}:{skill_dir}" (skill_dir = folder of the SKILL.md, "." at repo root).
Streaming: download one raw part -> convert with DuckDB -> upload -> delete. Resumable via a done-file.
Usage: python code/build_gitskills.py [artifacts|siblings|all]
"""
import os, sys, subprocess, time
from pathlib import Path
import duckdb
import pyarrow as pa, pyarrow.parquet as pq
sys.path.insert(0, str(Path(__file__).parent))
from lib_skill import cjk_ratio
from dotenv import load_dotenv
ROOT = Path(__file__).resolve().parent.parent
load_dotenv(ROOT / ".env")
BUCKET = "hf://buckets/Mercity/SkillsStorage"
RAW = f"{BUCKET}/raw/hf/mvaccargiu__gitskills/data"
OUT = f"{BUCKET}/processed/gitskills"
WORK = ROOT / "work/gs"; WORK.mkdir(parents=True, exist_ok=True)
DONE = ROOT / "logs/build_gitskills.done"
HF = str(ROOT / ".venv/bin/hf")
done = set(DONE.read_text().split()) if DONE.exists() else set()
def sh(*a):
for i in range(5):
r = subprocess.run(a, capture_output=True, text=True)
if r.returncode == 0: return
time.sleep(10 * (i + 1))
raise RuntimeError(f"{a}: {r.stderr[-500:]}")
def db():
c = duckdb.connect()
c.execute(f"SET memory_limit='650MB'; SET threads=1; SET temp_directory='{WORK}/duck_tmp'; SET preserve_insertion_order=false")
return c
def _unused_repos_table(c):
p = WORK / "repos.parquet"
if not p.exists(): sh(HF, "buckets", "cp", f"{RAW}/repos/part-00000.parquet", str(p))
c.execute(f"create or replace table repos as select full_name, license, stars, forks, is_fork, language, created_at, pushed_at from '{p}'")
SK_SCHEMA = pa.schema([("skill_uid", pa.string()), ("source", pa.string()), ("repo", pa.string()), ("skill_dir", pa.string()),
("skill_md_path", pa.string()), ("filename", pa.string()), ("is_skill_md", pa.bool_()), ("location_class", pa.string()),
("file_sha", pa.string()), ("skill_md", pa.string()), ("name", pa.string()), ("description", pa.string()),
("frontmatter_valid", pa.bool_()), ("body_chars", pa.int64()), ("has_scripts", pa.bool_()), ("has_references", pa.bool_()),
("sibling_count", pa.int64()), ("sibling_bytes", pa.int64()), ("composition_truncated", pa.bool_()),
("first_commit_at", pa.string()), ("last_commit_at", pa.string()), ("commit_count", pa.int64()), ("cjk_ratio", pa.float64()),
("repo_license", pa.string()), ("repo_stars", pa.int64()), ("repo_forks", pa.int64()), ("repo_is_fork", pa.bool_()),
("repo_language", pa.string()), ("repo_created_at", pa.string()), ("repo_pushed_at", pa.string()), ("snapshot", pa.string())])
FI_SCHEMA = pa.schema([("skill_uid", pa.string()), ("source", pa.string()), ("repo", pa.string()), ("skill_dir", pa.string()),
("rel_path", pa.string()), ("size", pa.int64()), ("blob_sha", pa.string()), ("content", pa.string()),
("content_bin", pa.binary()), ("has_content", pa.bool_()), ("skipped", pa.string())])
_repos = None
def repos_meta():
"""full_name -> (license, stars, forks, is_fork, language, created_at, pushed_at)"""
global _repos
if _repos is None:
p = WORK / "repos.parquet"
if not p.exists(): sh(HF, "buckets", "cp", f"{RAW}/repos/part-00000.parquet", str(p))
t = pq.read_table(p, columns=["full_name", "license", "stars", "forks", "is_fork", "language", "created_at", "pushed_at"])
_repos = {r[0]: r[1:] for r in zip(*[t.column(i).to_pylist() for i in range(t.num_columns)])}
return _repos
def stream(path, columns, batch=200):
"""Bounded-memory batches even when a part is ONE huge row group (later GitSkills parts hold ~0.5 GB of text):
buffered stream reads + no pre-buffering decode page by page instead of materialising the column chunk."""
pf = pq.ParquetFile(path, buffer_size=1 << 20, pre_buffer=False)
yield from pf.iter_batches(batch_size=batch, columns=columns, use_threads=False)
class Out:
"""Accumulate rows, write row groups of <= ~4000 rows / 64 MB; atomic rename on close."""
def __init__(self, path, schema): self.path, self.schema, self.rows, self.bytes, self.w, self.n = path, schema, [], 0, None, 0
def add(self, row, nbytes):
self.rows.append(row); self.bytes += nbytes; self.n += 1
if len(self.rows) >= 4000 or self.bytes > 32 << 20: self.flush()
def flush(self):
if not self.rows: return
self.w = self.w or pq.ParquetWriter(str(self.path) + ".tmp", self.schema, compression="zstd")
self.w.write_table(pa.Table.from_pylist(self.rows, self.schema)); self.rows, self.bytes = [], 0
def close(self):
self.flush()
if self.w: self.w.close()
else: pq.write_table(self.schema.empty_table(), str(self.path) + ".tmp")
os.replace(str(self.path) + ".tmp", self.path)
def sdir(path): return path.rsplit("/", 1)[0] if "/" in path else "."
def b(v): return None if v is None else bool(v)
def artifacts(part):
src, occ, sk = WORK / f"a-{part}", WORK / f"occ-{part}", WORK / f"sk-{part}"
if not src.exists(): sh(HF, "buckets", "cp", f"{RAW}/artifacts/{part}", str(src))
c = db() # occurrences: small columns only, DuckDB is fine
c.execute(f"""copy (select repo_full_name as repo, path, filename, location_class, file_sha, dedup_primary
from '{src}') to '{occ}' (format parquet, compression zstd)""")
c.close()
repos = repos_meta(); out = Out(sk, SK_SCHEMA)
cols = ["repo_full_name", "path", "filename", "location_class", "file_sha", "content", "frontmatter_valid", "name",
"description", "body_chars", "dedup_primary", "first_commit_at", "last_commit_at", "commit_count", "sibling_count",
"sibling_bytes", "has_scripts", "has_references", "composition_truncated"]
for bt in stream(src, cols):
for a in bt.to_pylist():
if a["dedup_primary"] != 1: continue
d = sdir(a["path"]); content = a["content"] or ""
r = repos.get(a["repo_full_name"]) or (None,) * 7
out.add({"skill_uid": f"gh:{a['repo_full_name']}:{d}", "source": "gitskills", "repo": a["repo_full_name"],
"skill_dir": d, "skill_md_path": a["path"], "filename": a["filename"], "is_skill_md": a["filename"] == "SKILL.md",
"location_class": a["location_class"], "file_sha": a["file_sha"], "skill_md": a["content"], "name": a["name"],
"description": a["description"], "frontmatter_valid": b(a["frontmatter_valid"]), "body_chars": a["body_chars"],
"has_scripts": b(a["has_scripts"]), "has_references": b(a["has_references"]), "sibling_count": a["sibling_count"],
"sibling_bytes": a["sibling_bytes"], "composition_truncated": b(a["composition_truncated"]),
"first_commit_at": a["first_commit_at"], "last_commit_at": a["last_commit_at"], "commit_count": a["commit_count"],
"cjk_ratio": cjk_ratio(content), "repo_license": r[0], "repo_stars": r[1], "repo_forks": r[2],
"repo_is_fork": b(r[3]), "repo_language": r[4], "repo_created_at": r[5], "repo_pushed_at": r[6],
"snapshot": "2026-07"}, len(content))
out.close()
sh(HF, "buckets", "cp", str(occ), f"{OUT}/occurrences/{part}")
sh(HF, "buckets", "cp", str(sk), f"{OUT}/skills/{part}")
for p in (src, occ, sk): p.unlink()
return out.n
def siblings(part):
src, fo = WORK / f"s-{part}", WORK / f"f-{part}"
if not src.exists(): sh(HF, "buckets", "cp", f"{RAW}/artifact_siblings/{part}", str(src))
out = Out(fo, FI_SCHEMA)
for bt in stream(src, ["repo_full_name", "artifact_path", "entry_name", "entry_type", "entry_size", "entry_sha", "content", "skipped_reason"]):
for s in bt.to_pylist():
if s["entry_type"] != "file": continue
d = sdir(s["artifact_path"]); c = s["content"]
out.add({"skill_uid": f"gh:{s['repo_full_name']}:{d}", "source": "gitskills", "repo": s["repo_full_name"], "skill_dir": d,
"rel_path": s["entry_name"], "size": s["entry_size"], "blob_sha": s["entry_sha"], "content": c, "content_bin": None,
"has_content": c is not None,
"skipped": None if c is not None else (s["skipped_reason"] or "not_fetched_by_gitskills")}, len(c or ""))
out.close()
sh(HF, "buckets", "cp", str(fo), f"{OUT}/files/{part}")
for p in (src, fo): p.unlink()
return out.n
if __name__ == "__main__":
what = sys.argv[1] if len(sys.argv) > 1 else "all"
jobs = []
if what in ("artifacts", "all"): jobs += [("artifacts", f"part-{i:05d}.parquet", artifacts) for i in range(31)]
if what in ("siblings", "all"): jobs += [("siblings", f"part-{i:05d}.parquet", siblings) for i in range(45)]
for kind, part, fn in jobs:
key = f"{kind}/{part}"
if key in done: continue
t = time.time(); n = fn(part)
with open(DONE, "a") as f: f.write(key + "\n")
print(f"{key}: {n} rows in {time.time() - t:.0f}s", flush=True)
print("done", flush=True)

Xet Storage Details

Size:
9.58 kB
·
Xet hash:
da5f9503e7efb1078b60eaf3eee08a287802b1c7dba46c30a80b694b76d0a807

Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.