Buckets:

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