024157529eaaa41e89d10d7db2750e6d13a00a62
max
  Wed Aug 26 03:01:39 2026 -0700
genark: add syncFtp to update the local NCBI genomes mirror

New subcommand that brings /hive/data/outside/ncbi/genomes/{GCA,GCF} up to date
with NCBI.  It fetches the current assembly_summary files for GenBank and RefSeq,
diffs them against the last synced copy to find assemblies that are new or have
changed, mirrors those with lftp using the exclude list from fetchLftp.sh, and
appends one row per affected assembly to changes.tsv so later steps know what
moved.  With -n it works out the whole change set and writes changes.tsv without
downloading any assembly data.

diff --git src/utils/genark/genark src/utils/genark/genark
index ecf45acd47f..3d5d1f349cf 100755
--- src/utils/genark/genark
+++ src/utils/genark/genark
@@ -1,331 +1,671 @@
 #!/usr/bin/env python3
 """genark - utilities to manage UCSC GenArk assembly hubs.
 
 This is a general, subcommand-based tool. Subcommands will be added over time.
 
 Subcommands:
   addContrib <name>   Install a contributed track collection into the GenArk
                       assembly hubs. <name> is a subdirectory of
                       /hive/data/genomes/asmHubs/contrib/ that is laid out as
                       one directory per assembly accession (GCA_*/GCF_*) plus a
                       shared docs/ directory, e.g.:
 
                         contrib/<name>/GCA_000000000.0/trackDb.txt
                         contrib/<name>/GCA_000000000.0/*.bb  (and/or *.bw)
                         contrib/<name>/docs/*.html
 
                       For each accession it writes into that assembly's GenArk
                       *build directory* (asmHubs/{genbankBuild,refseqBuild}/...)
                       -- never the served /gbdb/genark or asmHubs/<acc> symlink
                       trees, which the build system regenerates. There it:
                         - creates <buildDir>/contrib/<name>/ with symlinks to
                           the collection's data files (.bb/.bw) and doc pages;
                         - writes a per-assembly <name>.trackDb.txt whose
                           bigDataUrl/html paths are rewritten to be hub-root
                           relative (contrib/<name>/...);
                         - wires that trackDb block into the assembly's useOneFile
                           hub.txt, between BEGIN/END markers (idempotent).
 
                       --remove uninstalls: strips the hub.txt block and removes
                       the contrib/<name>/ symlink dir.
 
                       NOTE: a full GenArk hub rebuild regenerates hub.txt, so
                       re-run addContrib afterwards (or add <name> to the build's
                       asmHubTrackDb.sh for a durable inclusion).
 
   checkContrib <name> [accessions...]
                       Run hubCheck on assembly hubs that carry the collection
                       (a random --sample N by default, or --all, or the listed
                       accessions) and report problems, separating contrib-specific
                       issues from the assemblies' own pre-existing hub warnings.
                       --noTracks does a faster structure-only check.
+
+  syncFtp             Update the local NCBI assembly mirror
+                      /hive/data/outside/ncbi/genomes/{GCA,GCF}/ against NCBI.
+                      Fetches the current assembly_summary_{genbank,refseq}.txt,
+                      diffs them against the last-synced copy to find new and
+                      changed assemblies, mirrors those with lftp (reusing the
+                      curated exclude list), and appends one row per affected
+                      assembly to /hive/data/outside/ncbi/genomes/changes.tsv.
+                      -n computes the full change set and writes changes.tsv
+                      without downloading any assembly data.
 """
 
 import argparse
+import concurrent.futures
+import datetime
 import glob
 import os
 import random
 import re
 import shutil
 import subprocess
 import sys
+import tempfile
+import time
 
 ASMHUBS = "/hive/data/genomes/asmHubs"
 CONTRIB = os.path.join(ASMHUBS, "contrib")
 
 ACC_RE = re.compile(r"^GC[AF]_[0-9]{9}\.[0-9]+$")
 
+# --- syncFtp: local NCBI genomes mirror --------------------------------------
+
+GENOMES = "/hive/data/outside/ncbi/genomes"
+FTP_HOST = "https://ftp.ncbi.nlm.nih.gov"
+# assembly dirs live under genomes/all/<GCA|GCF>/... on the server; the local
+# mirror drops the "all" component, so genomes/all/GCA/000/001/405/GCA_...  maps
+# to <dest>/GCA/000/001/405/GCA_...
+ALL_MARKER = "/genomes/all/"
+
+# per-accession-type authority file: the current "latest" assemblies
+SUMMARY_PATH = {
+    "GCA": "/genomes/genbank/assembly_summary_genbank.txt",
+    "GCF": "/genomes/refseq/assembly_summary_refseq.txt",
+}
+SUMMARY_NAME = {
+    "GCA": "assembly_summary_genbank.txt",
+    "GCF": "assembly_summary_refseq.txt",
+}
+
+# lftp mirror exclude patterns (regex, passed as -x), copied verbatim from
+# /hive/data/outside/ncbi/genomes/fetchLftp.sh -- the files NCBI ships that the
+# GenArk mirror does not keep.
+LFTP_EXCLUDES = [
+    "suppressed",
+    "Annotation_comparison",
+    "RefSeq_transcripts_alignments",
+    "RNASeq_coverage_graphs",
+    r".*_ani_contam_ranges\.tsv",
+    r".*_ani_report\.txt",
+    r".*_fcs_report\.txt",
+    r".*_gene_ontology\.gaf\.gz",
+    r".*_genomic\.gtf\.gz",
+    r".*_protein\.gpff\.gz",
+    r".*_translated_cds\.faa\.gz",
+    r".*_wgsmaster\.gbff\.gz",
+    r"annotation_hashes\.txt",
+    r"md5checksums\.txt",
+    r"uncompressed_checksums\.txt",
+]
+
+# generous per-assembly backstop; lftp's own net:timeout handles stalls
+PER_ASM_TIMEOUT = 4 * 3600
+CHANGES_HEADER = "#date\taccession\tgcx\tstatus\tlocalDir\tnFiles\tfiles\n"
+
 
 def accPath(acc):
     """3-3-3 hashed subpath for an accession, e.g.
     GCA_041900255.1 -> GCA/041/900/255/GCA_041900255.1"""
     d = acc[4:]                 # 041900255.1
     return os.path.join(acc[0:3], d[0:3], d[3:6], d[6:9], acc)
 
 
 def buildDir(acc):
     """Resolve the GenArk *build* directory for an accession, or None.
 
     GCA_018506965.2 ->
       /hive/data/genomes/asmHubs/genbankBuild/GCA/018/506/965/GCA_018506965.2_HG005_mat_hprc_f2
 
     This is the directory the GenArk build system owns and builds from. The
     served copies -- /gbdb/genark/<acc>/ and asmHubs/<acc>/ -- are nothing but
     symlinks pointing back here, and they are regenerated by the build system,
     so contrib data must NEVER be written into them; it goes here in the build
     directory. GCF_* assemblies build under refseqBuild, GCA_* under
     genbankBuild, and the on-disk directory name carries an assembly-name suffix
     (…_HG005_mat_hprc_f2) beyond the bare accession, so glob to find it."""
     subtree = "refseqBuild" if acc.startswith("GCF_") else "genbankBuild"
     stem = os.path.join(ASMHUBS, subtree, accPath(acc))
     for m in sorted(glob.glob(stem + "_*")) + [stem]:
         if os.path.isdir(m):
             return m
     return None
 
 
 HUBS_URL = "https://hgdownload.soe.ucsc.edu/hubs"
 
 
 def hubUrl(acc):
     """Public served hub.txt URL for a GenArk assembly accession."""
     return "%s/%s/hub.txt" % (HUBS_URL, accPath(acc))
 
 
 def symlink(target, linkPath, dryRun):
     """Create/replace an absolute symlink linkPath -> target."""
     if dryRun:
         print("  ln -sf %s %s" % (target, linkPath))
         return
     if os.path.islink(linkPath) or os.path.exists(linkPath):
         os.remove(linkPath)
     os.symlink(target, linkPath)
 
 
 def rewriteTrackDb(srcPath, name):
     """Return the trackDb text with local bigDataUrl/html paths made hub-root
     relative (contrib/<name>/...). Remote (http/https/ftp) bigDataUrls and
     already-prefixed paths are left untouched."""
     out = []
     prefix = "contrib/%s/" % name
     for line in open(srcPath):
         stripped = line.lstrip()
         indent = line[:len(line) - len(stripped)]
         m = re.match(r"(bigDataUrl|linkDataUrl)\s+(\S+)\s*$", stripped)
         if m:
             key, val = m.group(1), m.group(2)
             if not re.match(r"[a-z]+://", val) and not val.startswith(prefix):
                 val = prefix + os.path.basename(val)
             out.append("%s%s %s\n" % (indent, key, val))
             continue
         m = re.match(r"html\s+(\S+)\s*$", stripped)
         if m:
             val = m.group(1)
             if not val.startswith(prefix):
                 val = prefix + os.path.basename(val)
             out.append("%shtml %s\n" % (indent, val))
             continue
         out.append(line)
     return "".join(out)
 
 
 def wireHubTxt(hubTxt, name, block, remove, dryRun):
     """Insert/replace (or remove) the marked contrib block in the assembly's
     useOneFile hub.txt. Edits the real file the hub.txt symlink points at."""
     begin = "# BEGIN genark contrib: %s" % name
     end = "# END genark contrib: %s" % name
     blockRe = re.compile(
         r"\n*" + re.escape(begin) + r".*?" + re.escape(end) + r"\n?", re.DOTALL)
     realHub = os.path.realpath(hubTxt)
     text = open(realHub).read()
     newText = blockRe.sub("\n", text).rstrip("\n") + "\n"
     if not remove:
         newText += "\n%s\n%s\n%s\n" % (begin, block.rstrip("\n"), end)
     if dryRun:
         print("  %s %s" % ("unwire hub.txt:" if remove else "wire hub.txt:", realHub))
     else:
         with open(realHub, "w") as fh:
             fh.write(newText)
 
 
 def addContrib(args):
     """Install a contrib track collection into the GenArk assembly hubs:
     symlink its data files + docs into <buildDir>/contrib/<name>/, write a
     per-assembly <name>.trackDb.txt with hub-root-relative paths, and wire that
     block into each assembly's served hub.txt. Everything is written into the
     assembly's GenArk build directory (see buildDir); the served /gbdb/genark
     and asmHubs/<acc> symlink trees are left alone. --remove undoes all of it."""
     name = args.name.rstrip("/")
     root = os.path.join(CONTRIB, name)
     if not os.path.isdir(root):
         sys.exit("error: no such contrib collection: %s" % root)
     docsDir = os.path.join(root, "docs")
     docs = []
     if os.path.isdir(docsDir):
         docs = sorted(f for f in os.listdir(docsDir) if f.endswith(".html"))
 
     accs = sorted(d for d in os.listdir(root) if ACC_RE.match(d)
                   and os.path.isdir(os.path.join(root, d)))
     if not accs:
         sys.exit("error: no accession directories (GCA_*/GCF_*) under %s" % root)
 
     done = 0
     skipped = 0
     for acc in accs:
         accDir = os.path.join(root, acc)
         asmDir = buildDir(acc)
         if asmDir is None:
             sys.stderr.write("skip %s: no GenArk build directory under "
                              "%s/{genbankBuild,refseqBuild}\n" % (acc, ASMHUBS))
             skipped += 1
             continue
         dest = os.path.join(asmDir, "contrib", name)
         # the served single-file hub in the build dir (asmHubs/<acc>/hub.txt and
         # /gbdb/genark/<acc>/hub.txt are symlinks to this); asmId is the build
         # directory basename, which carries the assembly-name suffix.
         asmId = os.path.basename(asmDir)
         hubTxt = os.path.join(asmDir, "%s.singleFile.hub.txt" % asmId)
 
         if args.remove:
             if os.path.exists(hubTxt):
                 wireHubTxt(hubTxt, name, "", remove=True, dryRun=args.dry_run)
             if args.dry_run:
                 print("  rm -rf %s" % dest)
             elif os.path.isdir(dest):
                 shutil.rmtree(dest)
             done += 1
             continue
 
         if args.dry_run:
             print("# %s -> %s" % (acc, dest))
         else:
             os.makedirs(dest, exist_ok=True)
 
         # symlink data files (.bb / .bw) from the collection's accession dir
         for f in sorted(os.listdir(accDir)):
             if f.endswith((".bb", ".bw")):
                 symlink(os.path.join(accDir, f), os.path.join(dest, f), args.dry_run)
         # symlink shared doc pages (flat, so contrib/<name>/<doc> resolves)
         for d in docs:
             symlink(os.path.join(docsDir, d), os.path.join(dest, d), args.dry_run)
 
         # per-assembly trackDb with hub-root-relative paths, and wire it into hub.txt
         srcTdb = os.path.join(accDir, "trackDb.txt")
         if os.path.isfile(srcTdb):
             tdb = rewriteTrackDb(srcTdb, name)
             destTdb = os.path.join(dest, "%s.trackDb.txt" % name)
             if args.dry_run:
                 print("  write %s (%d bytes)" % (destTdb, len(tdb)))
             else:
                 with open(destTdb, "w") as fh:
                     fh.write(tdb)
             if os.path.exists(hubTxt):
                 wireHubTxt(hubTxt, name, tdb, remove=False, dryRun=args.dry_run)
             else:
                 sys.stderr.write("warn %s: no hub.txt to wire at %s\n" % (acc, hubTxt))
         done += 1
 
     verb = "removed from" if args.remove else "installed into"
     print("addContrib %s: %s %d assemblies, skipped %d (no assembly hub)"
           % (name, verb, done, skipped))
 
 
 def contribTrackNames(root, accs):
     """Track names defined by the collection (from any one accession trackDb)."""
     for acc in accs:
         tdb = os.path.join(root, acc, "trackDb.txt")
         if os.path.isfile(tdb):
             return set(re.findall(r"^track\s+(\S+)", open(tdb).read(), re.MULTILINE))
     return set()
 
 
 def checkContrib(args):
     """Run hubCheck on a set of assembly hubs that carry the collection, and
     classify any reported problems as contrib-specific vs pre-existing hub
     warnings (so the collection can be signed off without wading through the
     assemblies' own tracks)."""
     name = args.name.rstrip("/")
     root = os.path.join(CONTRIB, name)
     if not os.path.isdir(root):
         sys.exit("error: no such contrib collection: %s" % root)
     accs = sorted(d for d in os.listdir(root) if ACC_RE.match(d)
                   and os.path.isdir(os.path.join(root, d)))
     tracks = contribTrackNames(root, accs)
 
     if args.accession:
         sel = args.accession
     elif args.all:
         sel = accs
     else:
         n = min(args.sample, len(accs))
         sel = sorted(random.sample(accs, n))
 
     cmd = ["hubCheck"]
     if args.noTracks:
         cmd.append("-noTracks")
 
     clean = other = flagged = failed = 0
     for acc in sel:
         try:
             res = subprocess.run(cmd + [hubUrl(acc)], capture_output=True,
                                  text=True, timeout=args.timeout)
             out = res.stdout + res.stderr
         except subprocess.TimeoutExpired:
             print("%s: TIMEOUT" % acc); failed += 1; continue
         problems = [l for l in out.splitlines()
                     if l.strip() and not l.startswith("Found ")]
         contribProblems = [l for l in problems
                            if any(t in l for t in tracks) or "rror" in l]
         if not problems:
             print("%s: clean" % acc); clean += 1
         elif contribProblems:
             print("%s: CONTRIB PROBLEMS (%d)" % (acc, len(contribProblems)))
             for l in contribProblems:
                 print("   ! %s" % l.strip())
             flagged += 1
         else:
             print("%s: ok (%d pre-existing hub warning(s), no contrib issue)"
                   % (acc, len(problems)))
             other += 1
 
     print("checkContrib %s: %d checked -- %d clean, %d ok-with-hub-warnings, "
           "%d CONTRIB PROBLEMS, %d failed"
           % (name, len(sel), clean, other, flagged, failed))
     if flagged or failed:
         sys.exit(1)
 
 
+def destDirFor(ftpPath, destBase):
+    """Map an assembly_summary ftp_path (column 20) to (localDir, srcSub), or
+    None. ftp_path e.g.
+      https://ftp.ncbi.nlm.nih.gov/genomes/all/GCA/000/001/405/GCA_000001405.28_GRCh38.p13
+    -> ("<destBase>/GCA/000/001/405/GCA_000001405.28_GRCh38.p13",
+        "GCA/000/001/405/GCA_000001405.28_GRCh38.p13").  srcSub is the path under
+    genomes/all/ used to build the remote lftp source."""
+    i = ftpPath.find(ALL_MARKER)
+    if i < 0:
+        return None
+    sub = ftpPath[i + len(ALL_MARKER):].rstrip("/")
+    if not sub:
+        return None
+    return os.path.join(destBase, sub), sub
+
+
+def fetchSummary(url, dest):
+    """Download url to dest with wget. Raises on failure."""
+    subprocess.run(["wget", "-q", "--tries=3", "--timeout=1200", "-O", dest, url],
+                   check=True)
+
+
+def summaryIndex(path):
+    """Stream a summary file into a compact {accession: hash(row)} dict. Only the
+    per-row hash is kept, so the 1.7 GB GenBank summary does not sit in RAM as
+    text."""
+    idx = {}
+    with open(path) as fh:
+        for line in fh:
+            if line.startswith("#"):
+                continue
+            acc = line.split("\t", 1)[0]
+            if acc:
+                idx[acc] = hash(line)
+    return idx
+
+
+def diffSummary(baselinePath, nowPath, gcx, destBase, fillMissing):
+    """Compare the last-synced baseline against the freshly fetched summary and
+    return (candidates, removed).  Both streamed; only a compact baseline index
+    is held in memory.
+      candidates: dicts {acc, gcx, status(new|updated), ftpPath, destDir, sub}
+      removed:    dicts {acc, gcx, ftpPath, destDir}  (dropped from "latest")
+    With fillMissing, unchanged rows whose local dir is absent are also
+    re-fetched (status new) to self-heal gaps from past partial runs."""
+    prev = summaryIndex(baselinePath)
+    candidates = []
+    with open(nowPath) as fh:
+        for line in fh:
+            if line.startswith("#"):
+                continue
+            f = line.rstrip("\n").split("\t")
+            if len(f) < 20:
+                continue
+            acc, ftpPath = f[0], f[19]
+            if not acc or not ftpPath.startswith("http"):
+                continue
+            dd = destDirFor(ftpPath, destBase)
+            if dd is None:
+                continue
+            destDir, sub = dd
+            if acc not in prev:
+                status = "new"
+            elif prev.pop(acc) == hash(line):
+                if fillMissing and not os.path.isdir(destDir):
+                    status = "new"
+                else:
+                    continue
+            else:
+                status = "updated"
+            candidates.append({"acc": acc, "gcx": gcx, "status": status,
+                               "ftpPath": ftpPath, "destDir": destDir, "sub": sub})
+
+    removed = []
+    if prev:                            # accessions left in prev dropped to historical
+        leftover = set(prev)
+        with open(baselinePath) as fh:
+            for line in fh:
+                if line.startswith("#") or not leftover:
+                    continue
+                f = line.rstrip("\n").split("\t")
+                if len(f) < 20 or f[0] not in leftover:
+                    continue
+                dd = destDirFor(f[19], destBase)
+                removed.append({"acc": f[0], "gcx": gcx, "ftpPath": f[19],
+                                "destDir": dd[0] if dd else ""})
+                leftover.discard(f[0])
+    return candidates, removed
+
+
+def parseLftpLog(logPath, host, sub, destDir):
+    """Read an lftp mirror --log/--script file and return (added, removed): the
+    file paths (relative to the assembly dir) that were, or would be,
+    transferred (get) and deleted (rm)."""
+    remotePrefix = "%s%s%s/" % (host, ALL_MARKER, sub)
+    added, removed = [], []
+    if not os.path.exists(logPath):
+        return added, removed
+    with open(logPath) as fh:
+        for line in fh:
+            parts = line.split()
+            if not parts:
+                continue
+            if parts[0] == "get":
+                url = parts[-1]
+                rel = url[len(remotePrefix):] if url.startswith(remotePrefix) \
+                    else os.path.basename(url)
+                added.append(rel)
+            elif parts[0] == "rm":
+                local = parts[-1]
+                removed.append(os.path.relpath(local, destDir))
+    return added, removed
+
+
+def lftpMirrorOne(cand, host, noDownload):
+    """Mirror one assembly with lftp (reusing the curated exclude list). Returns
+    the candidate dict augmented with added[], removed[], rc, err. With
+    noDownload, lftp runs in --script mode: it lists the remote dir and records
+    what it would transfer, but downloads nothing."""
+    destDir, sub = cand["destDir"], cand["sub"]
+    fd, logPath = tempfile.mkstemp(prefix="genarkSyncFtp.", suffix=".lftp")
+    os.close(fd)
+    logOpt = "--script=%s" % logPath if noDownload else "--log=%s" % logPath
+    excl = " ".join("-x '%s'" % p for p in LFTP_EXCLUDES)
+    script = ("open %s; set net:timeout 1200; "
+              "mirror --only-newer --no-perms --delete --parallel=4 %s %s "
+              "/genomes/all/%s/ %s/; quit"
+              % (host, logOpt, excl, sub, destDir))
+    if not noDownload:
+        os.makedirs(destDir, exist_ok=True)
+
+    rc, err = 1, ""
+    for attempt in range(3):
+        try:
+            res = subprocess.run(["lftp", "-e", script], capture_output=True,
+                                 text=True, timeout=PER_ASM_TIMEOUT)
+            rc, err = res.returncode, res.stderr
+            if rc == 0:
+                break
+        except subprocess.TimeoutExpired:
+            rc, err = 1, "timeout"
+        time.sleep(5 * (attempt + 1))
+
+    added, removed = parseLftpLog(logPath, host, sub, destDir)
+    os.remove(logPath)
+    cand = dict(cand)
+    cand.update(added=added, removed=removed, rc=rc, err=err)
+    return cand
+
+
+def writeChange(fh, stamp, rec):
+    """Append one assembly's change row to changes.tsv (fh)."""
+    files = list(rec.get("added", []))
+    files += ["rm:%s" % r for r in rec.get("removed", [])]
+    fh.write("%s\t%s\t%s\t%s\t%s\t%d\t%s\n"
+             % (stamp, rec["acc"], rec["gcx"], rec["status"],
+                rec.get("destDir", ""), len(files), ",".join(files)))
+    fh.flush()
+
+
+def syncFtp(args):
+    """Update the local /hive/data/outside/ncbi/genomes/{GCA,GCF}/ mirror:
+    fetch the current assembly summaries, diff against the last-synced copy to
+    find new and changed assemblies, mirror those with lftp, and log one row per
+    affected assembly to changes.tsv."""
+    host = args.src
+    destBase = args.dest
+    syncDir = os.path.join(destBase, "syncFtp")
+    changesTsv = os.path.join(destBase, "changes.tsv")
+    if args.gca_only:
+        gcxList = ["GCA"]
+    elif args.gcf_only:
+        gcxList = ["GCF"]
+    else:
+        gcxList = ["GCA", "GCF"]
+
+    if args.dry_run:
+        # global --dry-run: coarse preview only, touch nothing
+        for gcx in gcxList:
+            print("# would fetch %s%s and sync %s deltas into %s/%s/"
+                  % (host, SUMMARY_PATH[gcx], gcx, destBase, gcx))
+        return
+
+    os.makedirs(syncDir, exist_ok=True)
+    totNew = totUpd = totRem = totFail = 0
+    stamp = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
+    changesFh = None
+
+    for gcx in gcxList:
+        cur = os.path.join(syncDir, SUMMARY_NAME[gcx])
+        staged = cur + ".new"
+        print("# %s: fetching %s%s" % (gcx, host, SUMMARY_PATH[gcx]))
+        fetchSummary(host + SUMMARY_PATH[gcx], staged)
+
+        if not os.path.exists(cur):
+            os.rename(staged, cur)      # bootstrap: baseline only, no downloads
+            n = sum(1 for l in open(cur) if not l.startswith("#"))
+            print("# %s: bootstrap baseline saved (%d assemblies); no downloads. "
+                  "Re-run to sync deltas." % (gcx, n))
+            continue
+
+        candidates, removed = diffSummary(cur, staged, gcx, destBase,
+                                          args.fill_missing)
+        if args.limit is not None:
+            candidates = candidates[:args.limit]
+        nNew = sum(1 for c in candidates if c["status"] == "new")
+        print("# %s: %d new, %d updated, %d removed to process%s"
+              % (gcx, nNew, len(candidates) - nNew, len(removed),
+                 " (-n, no downloads)" if args.no_download else ""))
+
+        if changesFh is None:
+            newFile = not os.path.exists(changesTsv) \
+                or os.path.getsize(changesTsv) == 0
+            changesFh = open(changesTsv, "a")
+            if newFile:
+                changesFh.write(CHANGES_HEADER)
+
+        for rec in removed:
+            rec = dict(rec, status="removed")
+            writeChange(changesFh, stamp, rec)
+            totRem += 1
+
+        with concurrent.futures.ThreadPoolExecutor(max_workers=args.jobs) as ex:
+            futs = [ex.submit(lftpMirrorOne, c, host, args.no_download)
+                    for c in candidates]
+            for fut in concurrent.futures.as_completed(futs):
+                rec = fut.result()
+                if rec["rc"] != 0:
+                    totFail += 1
+                    sys.stderr.write("fail %s: lftp rc=%d %s\n"
+                                     % (rec["acc"], rec["rc"],
+                                        rec["err"].strip()[:200]))
+                    continue
+                writeChange(changesFh, stamp, rec)
+                if rec["status"] == "new":
+                    totNew += 1
+                else:
+                    totUpd += 1
+
+        # advance the baseline only on a real, successful run; -n leaves the
+        # baseline untouched so a later real run re-detects the same deltas.
+        if args.no_download:
+            os.remove(staged)
+        else:
+            if os.path.exists(cur):
+                os.replace(cur, cur + ".prev")
+            os.rename(staged, cur)
+
+    if changesFh is not None:
+        changesFh.close()
+    print("syncFtp: %d new, %d updated, %d removed, %d failed%s"
+          % (totNew, totUpd, totRem, totFail,
+             "  (dry, -n)" if args.no_download else ""))
+    if totFail:
+        sys.exit(1)
+
+
 def main():
     ap = argparse.ArgumentParser(
         prog="genark", description="Manage UCSC GenArk assembly hubs.")
     ap.add_argument("--dry-run", action="store_true",
                     help="show what would be done without changing anything")
     sub = ap.add_subparsers(dest="cmd", required=True)
 
     p = sub.add_parser("addContrib",
                        help="install a contrib track collection into the assembly hubs")
     p.add_argument("name", help="contrib subdirectory name under %s" % CONTRIB)
     p.add_argument("--remove", action="store_true",
                    help="uninstall: remove the symlinks and the hub.txt block")
     p.set_defaults(func=addContrib)
 
     c = sub.add_parser("checkContrib",
                        help="run hubCheck on assembly hubs carrying the collection")
     c.add_argument("name", help="contrib subdirectory name under %s" % CONTRIB)
     c.add_argument("accession", nargs="*",
                    help="specific accessions to check (default: a random sample)")
     c.add_argument("--sample", type=int, default=5,
                    help="number of random assemblies to check (default 5)")
     c.add_argument("--all", action="store_true", help="check every assembly")
     c.add_argument("--noTracks", action="store_true",
                    help="structure only, do not fetch track data (faster)")
     c.add_argument("--timeout", type=int, default=300,
                    help="per-assembly hubCheck timeout in seconds (default 300)")
     c.set_defaults(func=checkContrib)
 
+    s = sub.add_parser("syncFtp",
+                       help="update the local NCBI genomes mirror from the FTP site")
+    g = s.add_mutually_exclusive_group()
+    g.add_argument("--gca-only", action="store_true",
+                   help="only sync GenBank (GCA) assemblies")
+    g.add_argument("--gcf-only", action="store_true",
+                   help="only sync RefSeq (GCF) assemblies")
+    s.add_argument("-n", "--no-download", action="store_true",
+                   help="compute the change set and write changes.tsv, but "
+                        "download no assembly data (differs from --dry-run, "
+                        "which touches nothing)")
+    s.add_argument("--jobs", type=int, default=6,
+                   help="assemblies to mirror in parallel (default 6, max 20)")
+    s.add_argument("--limit", type=int,
+                   help="process at most this many candidate assemblies")
+    s.add_argument("--fill-missing", action="store_true",
+                   help="also re-fetch summary assemblies whose local dir is absent")
+    s.add_argument("--src", default=FTP_HOST,
+                   help="source FTP/HTTP host (default %s)" % FTP_HOST)
+    s.add_argument("--dest", default=GENOMES,
+                   help="local mirror root (default %s)" % GENOMES)
+    s.set_defaults(func=syncFtp)
+
     args = ap.parse_args()
+    if getattr(args, "jobs", 1) > 20:
+        args.jobs = 20
     args.func(args)
 
 
 if __name__ == "__main__":
     main()