Buckets:
| """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.