Buckets:
| """Convert raw GitHub bundle shards (code/github_bundles.py output) into the unified processed layout: | |
| processed/github_delta/<run>/skills/<shard>.parquet one row per bundle: SKILL.md text + parsed fields | |
| processed/github_delta/<run>/files/<shard>.parquet every other file of the bundle | |
| skill_uid = "gh:{repo}:{skill_dir}@{commit[:12]}" (GitSkills uses "gh:{repo}:{skill_dir}" for its Jul-2026 snapshot). | |
| Shards are self-contained (a repo's rows are never split across shards), so each converts independently. | |
| Usage: python code/build_github_delta.py RAW_DIR RUN_NAME (one-shot, local output only) | |
| see code/ship_delta.py for the convert -> upload -> delete loop used during crawls | |
| """ | |
| import sys | |
| 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 | |
| def convert_shard(shard, out_dir, batch_rows=2000): | |
| """Streams the shard in record batches (bounded RAM). Returns (skills_path, files_path or None, n_skills, n_files). | |
| Content is present exactly when `skipped` is null (crawler invariant).""" | |
| (out_dir / "skills").mkdir(parents=True, exist_ok=True); (out_dir / "files").mkdir(parents=True, exist_ok=True) | |
| sk_p, fi_p = out_dir / "skills" / Path(shard).name, out_dir / "files" / Path(shard).name | |
| pf = pq.ParquetFile(shard) | |
| # pass 1 (metadata columns only): per-bundle file counts and bytes | |
| agg = {} | |
| for b in pf.iter_batches(batch_size=50000, columns=["repo", "commit", "skill_dir", "rel_path", "size", "skipped"]): | |
| for r in b.to_pylist(): | |
| if r["rel_path"] == "SKILL.md": continue | |
| a = agg.setdefault(f"gh:{r['repo']}:{r['skill_dir']}@{(r['commit'] or 'unknown')[:12]}", [0, 0, 0]) | |
| a[0] += 1; a[1] += r["size"] or 0; a[2] += r["skipped"] is None | |
| # pass 2: stream full rows | |
| sk_w = fi_w = None; n_sk = n_fi = 0 | |
| for b in pf.iter_batches(batch_size=batch_rows): | |
| skills, files = [], [] | |
| for r in b.to_pylist(): | |
| uid = f"gh:{r['repo']}:{r['skill_dir']}@{(r['commit'] or 'unknown')[:12]}" | |
| text = r["content"].decode("utf-8", "replace") if r["content"] is not None and r["is_text"] else None | |
| if r["rel_path"] == "SKILL.md": | |
| ok, name, desc, fm, body = parse_skill_md(text or "") | |
| n, nb, nc = agg.get(uid, (0, 0, 0)) | |
| skills.append({"skill_uid": uid, "source": "github_delta", "repo": r["repo"], "commit": r["commit"], | |
| "skill_dir": r["skill_dir"], "skill_md_path": r["skill_md_path"], "file_sha": r["blob_sha"], | |
| "skill_md": text, "name": name, "description": desc, "frontmatter_valid": ok, | |
| "body_chars": len(body or ""), "cjk_ratio": cjk_ratio(text or ""), | |
| "n_bundle_files": n, "n_bundle_files_with_content": nc, "bundle_bytes": nb, | |
| "fetched_at": r["fetched_at"], "via": r["via"]}) | |
| else: | |
| files.append({"skill_uid": uid, "source": "github_delta", "repo": r["repo"], "skill_dir": r["skill_dir"], | |
| "rel_path": r["rel_path"], "size": r["size"], "blob_sha": r["blob_sha"], "content": text, | |
| "content_bin": r["content"] if r["content"] is not None and not r["is_text"] else None, | |
| "has_content": r["content"] is not None, "skipped": r["skipped"]}) | |
| if skills: | |
| t = pa.Table.from_pylist(skills, SK_SCHEMA); sk_w = sk_w or pq.ParquetWriter(sk_p, SK_SCHEMA, compression="zstd") | |
| sk_w.write_table(t); n_sk += len(skills) | |
| if files: | |
| t = pa.Table.from_pylist(files, FI_SCHEMA); fi_w = fi_w or pq.ParquetWriter(fi_p, FI_SCHEMA, compression="zstd") | |
| fi_w.write_table(t); n_fi += len(files) | |
| if sk_w: sk_w.close() | |
| else: pq.write_table(SK_SCHEMA.empty_table(), sk_p) | |
| if fi_w: fi_w.close() | |
| return sk_p, (fi_p if n_fi else None), n_sk, n_fi | |
| SK_SCHEMA = pa.schema([("skill_uid", pa.string()), ("source", pa.string()), ("repo", pa.string()), ("commit", pa.string()), | |
| ("skill_dir", pa.string()), ("skill_md_path", pa.string()), ("file_sha", pa.string()), ("skill_md", pa.string()), | |
| ("name", pa.string()), ("description", pa.string()), ("frontmatter_valid", pa.bool_()), ("body_chars", pa.int64()), | |
| ("cjk_ratio", pa.float64()), ("n_bundle_files", pa.int64()), ("n_bundle_files_with_content", pa.int64()), | |
| ("bundle_bytes", pa.int64()), ("fetched_at", pa.string()), ("via", 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())]) | |
| if __name__ == "__main__": | |
| if sys.argv[1] == "--one": # used by ship_delta.py: one shard per short-lived process | |
| _, _, n_sk, n_fi = convert_shard(Path(sys.argv[2]), Path(sys.argv[3])); print(n_sk, n_fi); sys.exit(0) | |
| raw, run = Path(sys.argv[1]), sys.argv[2] | |
| out = ROOT / "data/processed/github_delta" / run | |
| for shard in sorted(raw.glob("files-*.parquet")): | |
| print(shard.name, convert_shard(shard, out)[2:], flush=True) | |
Xet Storage Details
- Size:
- 5.45 kB
- Xet hash:
- 231cf1e6b3260040179560ca38766f2d3a895c96af9699a5fd7ed0d5dd31ffb2
·
Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.