Buckets:

Mercity/SkillsStorage / code /github_bundles.py
Pranav2748's picture
download
raw
17 kB
"""Bucket B: fetch full skill bundles (SKILL.md + every sibling file) from GitHub repos.
Per repo, two fetch paths:
git blobless shallow clone (tree only) -> stream the tree to find SKILL.md files -> list only the
skill dirs -> batch-fetch only bundle blobs.
codeload tarball of HEAD (one GET), used when anonymous git transfers get throttled (HTTP 401 on
upload-pack); streamed in passes so memory stays bounded. Git blob SHAs are recomputed so both
paths dedup against GitSkills file_sha.
The skill dir is the SKILL.md's folder; files of nested child skills belong to the child, not the parent.
A SKILL.md at the repo root makes the whole repo its bundle (listing capped at ROOT_LIST_CAP entries).
Memory is bounded per repo (huge monorepos never get their full file list held in RAM), shards are written
atomically with unique names, and repos are marked done only after their rows are on disk — so a separate
uploader can ship + delete finished shards while this runs.
Usage: python code/github_bundles.py SEEDS.txt OUT_DIR [workers]
"""
import os, sys, json, shutil, subprocess, threading, time, hashlib, tarfile, tempfile, itertools
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path, PurePosixPath
import requests, ctypes
import pyarrow as pa, pyarrow.parquet as pq
SEEDS, OUT = Path(sys.argv[1]), Path(sys.argv[2])
WORKERS = int(sys.argv[3]) if len(sys.argv) > 3 else 6
WORK = Path("/root/skills-db/work/clones")
MAX_FILE = 1 << 20 # skip content of any single file > 1 MB (still listed)
MAX_BIN = 256 << 10 # skip content of binary files > 256 KB
MAX_SKILL_BYTES = 8 << 20 # per-bundle content cap
MAX_REPO_BYTES = 100 << 20 # per-repo content cap (beyond: skipped='repo_cap'; these are mirror repos)
BIG_REPO = 24 << 20 # repos above this read their blobs one at a time (RAM)
MAX_SKILL_FILES = 500 # per-bundle file cap (SKILL.md first, then shallow paths)
MAX_SKILLS_PER_REPO = 3000 # mirror repos beyond this are recorded as deferred
ROOT_LIST_CAP = 5000 # root-level SKILL.md: list at most this many repo files as its bundle
MAX_TARBALL = 300 << 20 # codeload path: give up on bigger repos (retry via git later)
SHARD_ROWS = 20000
SHARD_BYTES = 64 << 20 # also flush when buffered content exceeds this (3.8 GB RAM box)
ENV = {**os.environ, "GIT_TERMINAL_PROMPT": "0", "GIT_LFS_SKIP_SMUDGE": "1"}
OUT.mkdir(parents=True, exist_ok=True); WORK.mkdir(parents=True, exist_ok=True)
done_f = OUT / "done.txt"
done = set(done_f.read_text().split()) if done_f.exists() else set()
# repos whose tarball exceeded MAX_TARBALL before: only the git path can fetch them, never re-download the tarball
needs_git = set()
if (OUT / "repos.jsonl").exists():
for _l in open(OUT / "repos.jsonl"):
if "tarball_too_large" in _l: needs_git.add(json.loads(_l)["repo"])
lock = threading.Lock(); rows = []; pending = []; buffered = [0]; seq = itertools.count()
git_blocked_until = [0.0]; throttle_hits = []; backoff = [150, 0.0] # [current pause s, last block time]
big_sem = threading.Semaphore(1)
stats = {"git": 0, "codeload": 0}
SCHEMA = pa.schema([("repo", pa.string()), ("commit", pa.string()), ("skill_dir", pa.string()),
("skill_md_path", pa.string()), ("rel_path", pa.string()), ("size", pa.int64()), ("blob_sha", pa.string()),
("content", pa.binary()), ("is_text", pa.bool_()), ("skipped", pa.string()), ("fetched_at", pa.string()),
("via", pa.string())])
try: _libc = ctypes.CDLL("libc.so.6")
except OSError: _libc = None
HTTP = requests.Session(); HTTP.headers["User-Agent"] = "skills-db-research/0.1"
class Throttled(Exception): pass
class Gone(Exception): pass
def git(*a, cwd=None, timeout=600, inp=None, env=ENV):
return subprocess.run(["git", *a], cwd=cwd, input=inp, capture_output=True, env=env, timeout=timeout)
def git_stream(*a, cwd=None):
"""Yield NUL-separated records of a git command's stdout without buffering it all."""
p = subprocess.Popen(["git", *a], cwd=cwd, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, env=ENV)
rest = b""
for chunk in iter(lambda: p.stdout.read(1 << 16), b""):
parts = (rest + chunk).split(b"\0"); rest = parts.pop()
for x in parts:
if x: yield x.decode(errors="replace")
if rest: yield rest.decode(errors="replace")
p.wait()
def blob_sha(data):
return hashlib.sha1(b"blob %d\0" % len(data) + data).hexdigest()
def flush(force=False):
"""Write buffered rows (atomically, unique name), then mark their repos done."""
with lock:
if not force and len(rows) < SHARD_ROWS and len(pending) < 500 and buffered[0] < SHARD_BYTES: return
batch, metas = rows[:], pending[:]; rows.clear(); pending.clear(); buffered[0] = 0
if batch:
name = OUT / f"files-{int(time.time() * 1000)}-{next(seq):05d}.parquet"
pq.write_table(pa.Table.from_pylist(batch, SCHEMA), str(name) + ".tmp", compression="zstd")
os.replace(str(name) + ".tmp", name)
with open(OUT / "repos.jsonl", "a") as f:
f.writelines(json.dumps(m) + "\n" for m in metas)
with open(done_f, "a") as f:
f.writelines(m["repo"] + "\n" for m in metas if m.get("ok"))
del batch
if _libc: _libc.malloc_trim(0) # glibc keeps freed blob buffers otherwise; return them to the OS
def skill_dir_of(p, dirset):
"""Nearest enclosing skill dir of path p (deepest wins), or None."""
d = p.rpartition("/")[0] or "."
while True:
if d in dirset: return d
if d == ".": return None
d = d.rpartition("/")[0] or "."
def plan(paths, dirs):
out = {d: [] for d in dirs}; dirset = set(dirs)
for p in paths:
d = skill_dir_of(p, dirset)
if d is not None: out[d].append(p)
return out
def dirs_of(skill_mds):
return sorted({p.rpartition("/")[0] or "." for p in skill_mds})
def order(paths): # SKILL.md first, then shallow paths
return sorted(paths, key=lambda p: (PurePosixPath(p).name != "SKILL.md", p.count("/"), p))
def select(bundles, size_of):
"""Apply count and byte caps (per file, per bundle, per repo). Returns {path: skip_reason or None}.
Every SKILL.md is selected before any sibling so the repo cap never drops a skill's own text."""
dec, repo_tot = {}, 0
for sd, files in bundles.items():
p = next((f for f in files if PurePosixPath(f).name == "SKILL.md"), None)
if p is not None and (size_of(p) or 0) <= MAX_FILE: repo_tot += size_of(p) or 0
for sd, files in bundles.items():
tot = 0
for i, p in enumerate(order(files)):
s = size_of(p); md = PurePosixPath(p).name == "SKILL.md"
if i >= MAX_SKILL_FILES: dec[p] = "bundle_cap"
elif s is None: dec[p] = "fetch_failed"
elif s > MAX_FILE: dec[p] = "too_large"
elif md: dec[p] = None; tot += s
elif tot + s > MAX_SKILL_BYTES: dec[p] = "bundle_cap"
elif repo_tot + s > MAX_REPO_BYTES: dec[p] = "repo_cap"
else: dec[p] = None; tot += s; repo_tot += s
return dec
def via_git(repo, d):
r = git("clone", "-q", "--depth", "1", "--filter=blob:none", "--no-checkout", "--single-branch",
f"https://github.com/{repo}.git", str(d), timeout=300)
if r.returncode:
err = r.stderr.decode(errors="replace")
if "could not read Username" in err or "429" in err or "rate limit" in err.lower(): raise Throttled(err[:200])
if "not found" in err.lower(): raise Gone(err[:200])
raise RuntimeError(err.strip()[:200])
commit = git("rev-parse", "HEAD", cwd=d).stdout.decode().strip()
# pass 1: stream names only, keep SKILL.md paths (bounded memory on monorepos)
n_files, mds = 0, []
for p in git_stream("ls-tree", "-r", "-z", "--name-only", "HEAD", cwd=d):
n_files += 1
if p == "SKILL.md" or p.endswith("/SKILL.md"): mds.append(p)
dirs = dirs_of(mds)
if not dirs or len(dirs) > MAX_SKILLS_PER_REPO: return commit, n_files, {d_: [] for d_ in dirs}, {}, {}
# pass 2: list only the skill dirs (root skill: whole repo, capped)
sha_of = {}
def take(rec):
info, path = rec.split("\t", 1); mode, typ, sha = info.split()
if typ == "blob": sha_of[path] = sha
if "." in dirs:
root_n = 0; dirset = set(dirs)
for rec in git_stream("ls-tree", "-r", "-z", "HEAD", cwd=d):
path = rec.split("\t", 1)[1]
if skill_dir_of(path, dirset) == ".":
root_n += 1
if root_n > ROOT_LIST_CAP: continue
take(rec)
else:
for k in range(0, len(dirs), 200):
for rec in git_stream("--literal-pathspecs", "ls-tree", "-r", "-z", "HEAD", "--", *dirs[k:k + 200], cwd=d):
take(rec)
bundles = plan(sha_of, dirs)
cand = list(dict.fromkeys(sha_of[p] for files in bundles.values() for p in order(files)[:MAX_SKILL_FILES]))
no_lazy = {**ENV, "GIT_NO_LAZY_FETCH": "1"}
for k in range(0, len(cand), 5000): # batched; explicit OIDs ignore --filter size limits
git("-c", "fetch.negotiationAlgorithm=noop", "fetch", "-q", "--no-tags", "--no-write-fetch-head",
"--recurse-submodules=no", "--filter=blob:none", "origin", "--stdin", cwd=d, timeout=900,
inp=("\n".join(cand[k:k + 5000]) + "\n").encode())
sizes = {}
for line in git("cat-file", "--batch-check", cwd=d, env=no_lazy, inp=("\n".join(cand) + "\n").encode()).stdout.decode().splitlines():
h = line.split()
if len(h) == 3 and h[1] == "blob": sizes[h[0]] = int(h[2])
dec = select(bundles, lambda p: sizes.get(sha_of[p]))
want = list(dict.fromkeys(sha_of[p] for p, why in dec.items() if why is None))
blobs = {}
big = sum(sizes.get(w, 0) for w in want) > BIG_REPO
if big: big_sem.acquire()
try:
if want:
buf = git("cat-file", "--batch", cwd=d, env=no_lazy, timeout=900, inp=("\n".join(want) + "\n").encode()).stdout
i = 0
while i < len(buf):
nl = buf.index(b"\n", i); hdr = buf[i:nl].decode().split()
if len(hdr) < 3 or hdr[1] == "missing": i = nl + 1; continue
n = int(hdr[2]); blobs[hdr[0]] = buf[nl + 1: nl + 1 + n]; i = nl + 2 + n
del buf
finally:
if big: big_sem.release()
files = {p: (sizes.get(sha_of[p], 0), sha_of[p], blobs.get(sha_of[p]) if dec[p] is None else None, dec[p])
for fs in bundles.values() for p in fs}
return commit, n_files, bundles, files, dec
def via_codeload(repo):
"""Tarball path using GNU tar (C) for listing/extraction; Python only reads the pax header for the commit."""
r = HTTP.get(f"https://codeload.github.com/{repo}/tar.gz/HEAD", stream=True, timeout=120)
if r.status_code == 404: raise Gone("codeload 404")
if r.status_code in (403, 429): raise Throttled(f"codeload {r.status_code}")
r.raise_for_status()
tmpd = Path(tempfile.mkdtemp(dir=WORK))
try:
tgz = tmpd / "repo.tgz"; n = 0
with open(tgz, "wb") as f:
for chunk in r.iter_content(1 << 20):
n += len(chunk)
if n > MAX_TARBALL: raise RuntimeError("tarball_too_large")
f.write(chunk)
with tarfile.open(tgz, mode="r|gz") as tf: # first header only: git archive puts the commit id in pax 'comment'
first = tf.next(); commit = tf.pax_headers.get("comment")
top = first.name.split("/", 1)[0] if first else ""
lst = subprocess.run(["tar", "--quoting-style=literal", "-tzf", str(tgz)], capture_output=True, timeout=600)
names = [x for x in lst.stdout.decode(errors="replace").split("\n") if x and not x.endswith("/")]
rels = [x.partition("/")[2] for x in names if x.partition("/")[2]]
n_files = len(rels)
dirs = dirs_of([p for p in rels if p == "SKILL.md" or p.endswith("/SKILL.md")])
if not dirs or len(dirs) > MAX_SKILLS_PER_REPO: return commit, n_files, {d_: [] for d_ in dirs}, {}, {}
dirset, cand, root_n = set(dirs), [], 0
for p in rels:
sd = skill_dir_of(p, dirset)
if sd is None: continue
if sd == ".":
root_n += 1
if root_n > ROOT_LIST_CAP: continue
cand.append(p)
bundles = plan(cand, dirs)
keep = [p for fs in bundles.values() for p in order(fs)[:MAX_SKILL_FILES]]
ex = tmpd / "x"; ex.mkdir()
(tmpd / "list").write_bytes(b"\0".join(f"{top}/{p}".encode() for p in keep) + b"\0")
subprocess.run(["tar", "-xzf", str(tgz), "-C", str(ex), "--null", "-T", str(tmpd / "list"), "--no-same-owner"],
capture_output=True, timeout=600)
tgz.unlink()
base = ex / top
def size_of(p):
fp = base / p
if not os.path.lexists(fp): return None
return len(os.readlink(fp).encode()) if os.path.islink(fp) else fp.stat().st_size
sizes = {p: size_of(p) for p in keep}
dec = select(bundles, lambda p: sizes.get(p))
files = {}
for fs in bundles.values():
for p in fs:
why = dec.get(p, "bundle_cap"); data = None
if why is None:
fp = base / p
data = os.readlink(fp).encode() if os.path.islink(fp) else fp.read_bytes()
files[p] = (sizes.get(p) or 0, blob_sha(data) if data is not None else None, data, why)
return commit, n_files, bundles, files, dec
finally:
shutil.rmtree(tmpd, ignore_errors=True)
def process(repo):
if repo in done: return
d = WORK / hashlib.md5(repo.encode()).hexdigest(); shutil.rmtree(d, ignore_errors=True)
meta = {"repo": repo, "ok": False}
try:
try:
if time.time() < git_blocked_until[0]: raise Throttled("git cooling down")
commit, nfiles, bundles, files, dec = via_git(repo, d); meta["via"] = "git"
except Throttled as e:
# git answers 401 both for throttling and for deleted/private repos: codeload tells them apart
shutil.rmtree(d, ignore_errors=True)
if repo in needs_git: raise RuntimeError("tarball_too_large (deferred: needs git)")
commit, nfiles, bundles, files, dec = via_codeload(repo); meta["via"] = "codeload"
if "cooling down" not in str(e): # repo exists, so git really refused: pause git after 3 in 60 s
with lock:
now_t = time.time(); throttle_hits[:] = [x for x in throttle_hits if now_t - x < 60] + [now_t]
if len(throttle_hits) >= 3:
backoff[0] = min(backoff[0] * 2, 3600) if now_t - backoff[1] < backoff[0] + 600 else 300
backoff[1] = now_t; git_blocked_until[0] = now_t + backoff[0]; throttle_hits.clear()
print(f"git throttled: pausing git for {backoff[0]}s, using codeload", flush=True)
meta.update(commit=commit, n_files_repo=nfiles, n_skills=len(bundles))
if len(bundles) > MAX_SKILLS_PER_REPO: meta["deferred"] = "too_many_skills"
now = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()); out = []
for sd, fs in (bundles.items() if files else []):
for p in fs:
s, sha, c, why = files[p]
text = c is not None and b"\0" not in c[:8192]
if c is not None and not text and len(c) > MAX_BIN: c, why = None, "binary_large"
elif why is None and c is None: why = "fetch_failed"
out.append({"repo": repo, "commit": commit, "skill_dir": sd,
"skill_md_path": "SKILL.md" if sd == "." else f"{sd}/SKILL.md",
"rel_path": p if sd == "." else p[len(sd) + 1:], "size": s, "blob_sha": sha,
"content": c, "is_text": text, "skipped": why, "fetched_at": now, "via": meta["via"]})
with lock: rows.extend(out); stats[meta["via"]] += 1; buffered[0] += sum(len(r["content"] or b"") for r in out)
meta.update(ok=True, n_rows=len(out), n_skipped=sum(1 for r in out if r["skipped"]))
except Gone as e:
meta.update(ok=True, gone=str(e)[:100])
except Exception as e:
meta["err"] = f"{type(e).__name__}: {e}"[:200]
finally:
shutil.rmtree(d, ignore_errors=True)
with lock: pending.append(meta)
flush()
if __name__ == "__main__":
repos = list(dict.fromkeys(l.strip() for l in open(SEEDS) if l.strip() and l.strip() not in done))
print(f"{len(repos)} repos to process, {WORKERS} workers", flush=True); t = time.time()
with ThreadPoolExecutor(WORKERS) as ex:
for i, _ in enumerate(ex.map(process, repos), 1):
if i % 200 == 0: print(f"{i}/{len(repos)} {i / (time.time() - t):.1f} repos/s {stats}", flush=True)
flush(force=True); print("done", stats, flush=True)

Xet Storage Details

Size:
17 kB
·
Xet hash:
dab75cdf5b3462edfb01fb61025223a1abbd4358887dd0b645ee5f41377345b4

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