ec8d25d5bd0eac0a4f5172603e704106cbf6c66b max Sun Oct 4 06:37:38 2026 -0700 methbase2 otto: ssh to hgdownload as the invoking user, not as qateam, refs #34246 diff --git src/hg/utils/otto/methbase2/methbaseDownload src/hg/utils/otto/methbase2/methbaseDownload index 88bc53891f4..8f39723590a 100755 --- src/hg/utils/otto/methbase2/methbaseDownload +++ src/hg/utils/otto/methbase2/methbaseDownload @@ -1,538 +1,538 @@ #!/usr/bin/env python3 """Mirror the MethBase2 track hub from smithlab.usc.edu to local disk. hubClone cannot do this job. The hub's bigDataUrls are absolute and point into http://smithlab.usc.edu/methbase/data/, a tree that is not under the hub directory at all, and hubClone -download writes every data file into a single flat per-assembly directory named after the URL's basename. That would collapse the SRA project/sample/experiment hierarchy into ~20,000 files per assembly directory, and it only skips a file if it already exists, so it can never repair a truncated download or pick up a file that changed upstream. It is also single-threaded. This script keeps the remote directory layout, downloads in parallel, and re-fetches a file only when the remote Content-Length differs from what is on disk, so it works both for the initial mirror and for later updates. Layout under --outDir: hub/ the hub's own text files hub.txt genomes.txt /trackDb.txt /mb2__metadata.tsv /mb2__colors.json data/ the bigDataUrl tree with --stripPrefix removed ///.hmr.bb ... Every bigDataUrl in the mirrored trackDb files is rewritten so it no longer reaches back to smithlab.usc.edu. By default it becomes an absolute url under --baseUrl (https://hgdownload.soe.ucsc.edu/hubs/methbase/v3/data/...), which is what makes the hub loadable by url from anywhere. Pass --relativeUrls for the "../../data/..." form instead, which keeps the tree self-contained so it can be moved or served from a different host, at the cost of only working when hub/ and data/ stay siblings. """ import argparse import os import sys import threading import time from concurrent.futures import ThreadPoolExecutor from urllib.parse import urljoin, urlparse import requests # The faceted hub. Note this is NOT the hub at .../methylation/hub.txt, which is # the older superTrack-per-study layout. methbase_hub/ is an alias for the # latest methbase_hub_/, and those dated directories are rewritten in # place, so the date in the name buys no reproducibility - prefer the alias. DEFAULT_HUB_URL = "http://smithlab.usc.edu/trackdata/methylation/methbase_hub/hub.txt" # The mirror lives on hgdownload, served as https://hgdownload.soe.ucsc.edu/hubs/ -# and reached as qateam@hgdownload. Each refresh gets its own vN directory, so +# served from hgdownload. Each refresh gets its own vN directory, so # bump this (or pass -o) when v4 is cut. Do not update the directory that is # being served: hardlink it first (cp -al v3 v4), update the copy, then swap. DEFAULT_OUT_DIR = "/mirrordata/hubs/methbase/v3" # Where that directory is served from. Rewritten bigDataUrls point here, so the # mirrored trackDb works for anyone who loads it by url, not only for a copy # sitting next to the data. Keep this and DEFAULT_OUT_DIR pointing at the same # vN directory. DEFAULT_BASE_URL = "https://hgdownload.soe.ucsc.edu/hubs/methbase/v3" DEFAULT_STRIP_PREFIX = "http://smithlab.usc.edu/methbase/data/" # trackDb settings whose value is a big binary file that lives outside the hub BIG_DATA_TAGS = frozenset([ "bigDataUrl", "bigDataIndex", "barChartMatrixUrl", "barChartSampleUrl", "linkDataUrl", "frames", "summary", "searchTrix", ]) # trackDb settings whose value is a small text file. When the value is a # relative name the file belongs to the hub and is mirrored under hub/ and # always re-fetched; when it is an absolute url it is treated as a data file # instead (see scanTrackDb). HUB_TEXT_TAGS = frozenset([ "metaDataUrl", "colorSettingsUrl", "html", "groups", ]) # how long to wait for the server, in seconds: (connect, read) TIMEOUT = (30, 300) # thread-local requests.Session, so each worker keeps its own keep-alive # connection to the server instead of reconnecting for every file _local = threading.local() def session(): "return this thread's requests.Session, making it on first use" s = getattr(_local, "session", None) if s is None: s = requests.Session() s.headers["User-Agent"] = "UCSC-methbaseDownload/1.0" _local.session = s return s def fetchText(url, retries=3): "GET url and return the body as text, retrying on transient errors" lastErr = None for attempt in range(retries): try: r = session().get(url, timeout=TIMEOUT) r.raise_for_status() return r.text except Exception as e: lastErr = e time.sleep(2 ** attempt) raise IOError("cannot fetch %s: %s" % (url, lastErr)) def remoteSize(url, retries=3): """return the Content-Length of url, or None if the server does not say. Raises on a hard error, e.g. a 404, so a vanished file gets reported instead of silently skipped.""" lastErr = None for attempt in range(retries): try: r = session().head(url, timeout=TIMEOUT, allow_redirects=True) r.raise_for_status() size = r.headers.get("Content-Length") return int(size) if size is not None else None except Exception as e: lastErr = e time.sleep(2 ** attempt) raise IOError("HEAD failed for %s: %s" % (url, lastErr)) def writeTextAtomic(localPath, text): """write a text file via a temp file and a rename. Never write in place: the mirror is served to the public while it updates, so a reader must not be able to catch a half-written trackDb, and updating a hardlinked staging copy (cp -al v3 v4) must not write through into the directory still being served.""" os.makedirs(os.path.dirname(localPath), exist_ok=True) tmp = localPath + ".tmp" with open(tmp, "w") as fh: fh.write(text) os.replace(tmp, localPath) def humanSize(n): "short byte count, e.g. 21K or 64.1M, so log lines stay narrow" for unit in ("B", "K", "M", "G", "T"): if abs(n) < 1024 or unit == "T": return "%d%s" % (n, unit) if unit == "B" else "%.1f%s" % (n, unit) n /= 1024.0 def isAbsoluteUrl(val): "True if a trackDb setting names a remote file rather than a hub-local one" return val.startswith("http://") or val.startswith("https://") def parseStanzas(text): """yield (tag, value) pairs from a trackDb/hub text file. Stanza boundaries do not matter here - we only need every setting that names a file - so this is a flat scan. Continuation lines ending in a backslash are joined.""" pending = "" for line in text.splitlines(): line = pending + line.rstrip("\n") pending = "" if line.endswith("\\"): pending = line[:-1] continue line = line.strip() if not line or line.startswith("#"): continue parts = line.split(None, 1) if len(parts) != 2: continue yield parts[0], parts[1].strip() class Mirror(object): "collects the download list from the hub text files, then fetches the data" def __init__(self, args): self.args = args self.hubDir = os.path.join(args.outDir, "hub") self.dataDir = os.path.join(args.outDir, "data") # url -> local path, deduplicated: the same file can be referenced from # more than one assembly self.dataFiles = {} self.lock = threading.Lock() self.nDone = 0 self.nSkipped = 0 self.nDownloaded = 0 self.bytesDownloaded = 0 self.failures = [] self.skippedGenomes = [] # url -> assembly that first referenced it, so the log can name the # assemblies a run actually touched self.fileDb = {} # (path, oldSize or None, newSize) for every file this run wrote self.changes = [] # ---------------------------------------------------------------- hub text def mirrorText(self, url, localPath): "fetch a hub text file, write it under hub/, and return its content" text = fetchText(url, self.args.retries) writeTextAtomic(localPath, text) return text def dataPath(self, url): """map a bigDataUrl to a path under data/. Normally the URL starts with --stripPrefix and the rest is used as-is. Anything else is filed under data/_other// so a stray URL still lands somewhere sane instead of aborting the run.""" if url.startswith(self.args.stripPrefix): rel = url[len(self.args.stripPrefix):] else: p = urlparse(url) rel = os.path.join("_other", p.netloc, p.path.lstrip("/")) sys.stderr.write("warning: url outside --stripPrefix, storing as " "data/%s: %s\n" % (rel, url)) rel = rel.split("?")[0].lstrip("/") return os.path.join(self.dataDir, rel) def linkFor(self, target, outDir): """what a rewritten bigDataUrl should point at. By default an absolute url under --baseUrl, so the mirrored trackDb works for anyone who loads it from hgdownload rather than only from a copy sitting next to the data. --relativeUrls gives the "../../data/.." form instead, which keeps the tree self-contained and relocatable.""" rel = os.path.relpath(target, self.dataDir) if self.args.relativeUrls: return os.path.relpath(target, outDir) return "%s/data/%s" % (self.args.baseUrl.rstrip("/"), rel) def scanTrackDb(self, url, localPath, db, seen): """mirror one trackDb.txt, noting and relinking every file it names. The copy written to disk has each bigDataUrl replaced by a relative path to where that file will land under data/, which is what makes the mirror loadable on its own. Follows include lines; `seen` guards against an include cycle.""" if url in seen: return seen.add(url) text = fetchText(url, self.args.retries) outDir = os.path.dirname(localPath) os.makedirs(outDir, exist_ok=True) includes = [] out = [] for line in text.splitlines(True): parts = line.strip().split(None, 1) tag = parts[0] if parts else None val = parts[1].strip() if len(parts) == 2 else None # A hub-text setting given as an absolute url is not hub-local # at all - MethBase2 writes its description pages as # "html http://smithlab.usc.edu/methbase/data/.../x.html" - so # file it under data/ and relink it like any other data file. # Left alone it would both mangle the local path and leave the # mirrored trackDb reaching back to smithlab for the page. if val and (tag in BIG_DATA_TAGS or (tag in HUB_TEXT_TAGS and isAbsoluteUrl(val))): dataUrl = urljoin(url, val) target = self.dataPath(dataUrl) self.dataFiles.setdefault(dataUrl, target) self.fileDb.setdefault(dataUrl, db) indent = line[:len(line) - len(line.lstrip())] out.append("%s%s %s\n" % (indent, tag, self.linkFor(target, outDir))) continue # include and the small hub-local text files keep their # relative names, so they need no rewriting - only fetching if tag == "include" and val: includes.append(val) elif tag in HUB_TEXT_TAGS and val: try: self.mirrorText(urljoin(url, val), os.path.join(outDir, val)) except IOError as e: sys.stderr.write("warning: %s\n" % e) out.append(line) writeTextAtomic(localPath, "".join(out)) for val in includes: self.scanTrackDb(urljoin(url, val), os.path.join(outDir, val), db, seen) def scanHub(self): "mirror hub.txt, genomes.txt and every trackDb, filling self.dataFiles" hubUrl = self.args.hubUrl hubText = self.mirrorText(hubUrl, os.path.join(self.hubDir, "hub.txt")) genomesFile = None for tag, val in parseStanzas(hubText): if tag == "genomesFile": genomesFile = val if genomesFile is None: raise SystemExit("no genomesFile setting in %s" % hubUrl) genomesUrl = urljoin(hubUrl, genomesFile) genomesText = self.mirrorText(genomesUrl, os.path.join(self.hubDir, genomesFile)) # genomes.txt is stanzas of "genome " followed by "trackDb " genomes = [] db = None for tag, val in parseStanzas(genomesText): if tag == "genome": db = val elif tag == "trackDb" and db is not None: genomes.append((db, val)) db = None if self.args.assemblies: wanted = set(self.args.assemblies) missing = wanted - set(d for d, _ in genomes) if missing: raise SystemExit("assemblies not in this hub: %s" % ", ".join(sorted(missing))) genomes = [g for g in genomes if g[0] in wanted] print("hub has %d assemblies, reading their trackDb files" % len(genomes)) for db, tdbPath in genomes: before = len(self.dataFiles) try: self.scanTrackDb(urljoin(genomesUrl, tdbPath), os.path.join(self.hubDir, tdbPath), db, set()) except IOError as e: # An assembly listed in genomes.txt whose directory was never # published upstream must not sink the whole mirror - as of # 2026-09-15 the MethBase2 faceted hub declares aplCal1 but # serves no aplCal1/trackDb.txt. Skipped assemblies are counted # and make the exit status non-zero, so this is never silent. sys.stderr.write("warning: skipping assembly %s: %s\n" % (db, e)) self.skippedGenomes.append(db) continue print(" %-10s %6d data files (%d new)" % (db, len(self.dataFiles) - before, len(self.dataFiles) - before)) # ----------------------------------------------------------------- fetching def needsDownload(self, url, path): """(needed, oldSize) - oldSize is None when the file is new. oldSize is what the log reports the size difference against. No HEAD is sent when the file does not exist yet, which keeps the initial mirror down to one request per file.""" if not os.path.exists(path): return True, None localSize = os.path.getsize(path) if self.args.noHead: return False, localSize rSize = remoteSize(url, self.args.retries) if rSize is None: # server did not give a length, so we cannot compare: keep what we # have rather than re-downloading everything on every run return False, localSize return rSize != localSize, localSize def download(self, url, path): "stream url into path, via a .tmp file so a crash leaves no partial" os.makedirs(os.path.dirname(path), exist_ok=True) tmp = path + ".tmp" lastErr = None for attempt in range(self.args.retries): try: nBytes = 0 with session().get(url, timeout=TIMEOUT, stream=True) as r: r.raise_for_status() with open(tmp, "wb") as fh: for chunk in r.iter_content(chunk_size=1024 * 1024): fh.write(chunk) nBytes += len(chunk) expected = r.headers.get("Content-Length") if expected is not None and int(expected) != nBytes: raise IOError("short read: got %d of %s bytes" % (nBytes, expected)) os.replace(tmp, path) return nBytes except Exception as e: lastErr = e if os.path.exists(tmp): os.remove(tmp) time.sleep(2 ** attempt) raise IOError("download failed for %s: %s" % (url, lastErr)) def fetchOne(self, item): url, path = item try: needed, oldSize = self.needsDownload(url, path) if not needed: with self.lock: self.nSkipped += 1 self.tick() return if self.args.dryRun: print("would download %s -> %s" % (url, path)) with self.lock: self.nDownloaded += 1 self.tick() return nBytes = self.download(url, path) with self.lock: self.nDownloaded += 1 self.bytesDownloaded += nBytes self.changes.append((url, path, oldSize, nBytes)) self.tick() except Exception as e: with self.lock: self.failures.append((url, str(e))) self.tick() sys.stderr.write("ERROR %s: %s\n" % (url, e)) def tick(self): "progress line; caller must hold self.lock" self.nDone += 1 if self.nDone % self.args.progressEvery == 0 or self.nDone == self.total: elapsed = time.time() - self.startTime print("%d/%d done %d downloaded (%.1f GB) %d unchanged " "%d failed %.0fs" % (self.nDone, self.total, self.nDownloaded, self.bytesDownloaded / 1e9, self.nSkipped, len(self.failures), elapsed), flush=True) def fetchAll(self): items = sorted(self.dataFiles.items()) self.total = len(items) self.startTime = time.time() print("%d distinct data files, %d threads" % (self.total, self.args.threads)) with ThreadPoolExecutor(max_workers=self.args.threads) as pool: list(pool.map(self.fetchOne, items)) def writeLog(self): """append one short record of this run to download.log. Kept terse on purpose: a separator, the date, the assemblies this run actually changed, then one line per file written - the new size, or the old and new size and the delta for a file that was replaced.""" logPath = os.path.join(self.args.outDir, "download.log") dbs = sorted({self.fileDb.get(url, "?") for url, _, _, _ in self.changes}) with open(logPath, "a") as fh: fh.write("%s\n" % ("-" * 64)) fh.write("%s %s\n" % (time.strftime("%Y-%m-%d %H:%M:%S"), " ".join(dbs) if dbs else "no changes")) for url, path, oldSize, newSize in sorted(self.changes, key=lambda c: c[1]): rel = os.path.relpath(path, self.args.outDir) if oldSize is None: fh.write("new %s %s\n" % (rel, humanSize(newSize))) else: fh.write("upd %s %s -> %s (%+d B)\n" % (rel, humanSize(oldSize), humanSize(newSize), newSize - oldSize)) if self.skippedGenomes: fh.write("skipped assemblies: %s\n" % " ".join(self.skippedGenomes)) if self.failures: fh.write("failed: %d (see failures.txt)\n" % len(self.failures)) fh.write("%d written, %.1f GB, %d unchanged\n" % (self.nDownloaded, self.bytesDownloaded / 1e9, self.nSkipped)) def report(self): print("\ndownloaded %d files (%.1f GB), %d already up to date, %d failed" % (self.nDownloaded, self.bytesDownloaded / 1e9, self.nSkipped, len(self.failures))) retVal = 0 if self.skippedGenomes: print("SKIPPED %d assemblies listed in genomes.txt whose trackDb " "could not be read: %s" % (len(self.skippedGenomes), ", ".join(self.skippedGenomes))) retVal = 1 if self.failures: failFile = os.path.join(self.args.outDir, "failures.txt") with open(failFile, "w") as fh: for url, err in self.failures: fh.write("%s\t%s\n" % (url, err)) print("failed urls written to %s - rerun the script to retry them" % failFile) retVal = 1 return retVal def main(): parser = argparse.ArgumentParser( description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) parser.add_argument("--hubUrl", default=DEFAULT_HUB_URL, help="hub.txt url (default: %(default)s)") parser.add_argument("-o", "--outDir", default=DEFAULT_OUT_DIR, help="output directory (default: %(default)s)") parser.add_argument("-t", "--threads", type=int, default=8, help="parallel downloads (default: %(default)s)") parser.add_argument("--stripPrefix", default=DEFAULT_STRIP_PREFIX, help="url prefix removed to get the path under data/ " "(default: %(default)s)") parser.add_argument("--baseUrl", default=DEFAULT_BASE_URL, help="public url the mirror is served at; rewritten " "bigDataUrls point under its data/ " "(default: %(default)s)") parser.add_argument("--relativeUrls", action="store_true", help="rewrite bigDataUrls as relative paths " "(../../data/...) instead of urls under " "--baseUrl, making the tree relocatable") parser.add_argument("-a", "--assemblies", nargs="+", metavar="DB", help="only these assemblies, e.g. -a hg38 mm39") parser.add_argument("-n", "--dryRun", action="store_true", help="list what would be downloaded, fetch nothing") parser.add_argument("--noHead", action="store_true", help="do not HEAD existing files, assume they are " "current; faster but misses upstream changes") parser.add_argument("--retries", type=int, default=3, help="attempts per url (default: %(default)s)") parser.add_argument("--progressEvery", type=int, default=200, help="print progress every N files (default: %(default)s)") args = parser.parse_args() # line buffering, so a run redirected to a log file shows progress as it # happens instead of in one block at the end sys.stdout.reconfigure(line_buffering=True) os.makedirs(args.outDir, exist_ok=True) mirror = Mirror(args) mirror.scanHub() mirror.fetchAll() if not args.dryRun: mirror.writeLog() return mirror.report() if __name__ == "__main__": sys.exit(main())