← back to Kravet Sheet Sync 2026 04 20

kravet_ftp_mirror.py

168 lines

#!/usr/bin/env python3
"""
Mirror Kravet SFTP (ftp.kravet.com:22) to local disk.

Strategy:
- 8 worker threads, each with its own paramiko Transport (SFTP needs 1 session per thread)
- Atomic writes: download to <path>.tmp, rename on success
- Resume: skip files where local size == remote size
- Progress: JSONL log line per completed file + periodic summary
- Crash-safe: re-running continues from where it left off

Run:
  source /root/.kravet-sftp.env && \\
  python3 kravet_ftp_mirror.py --dest /root/kravet-ftp-mirror \\
                               --threads 8 --log /root/kravet-ftp-mirror.log
"""
import argparse, json, os, stat, sys, threading, time, traceback
from concurrent.futures import ThreadPoolExecutor, as_completed
from queue import Queue
import paramiko

# ------------------------------------------------------------------
def load_env(path):
    env = {}
    with open(path) as f:
        for line in f:
            line = line.strip()
            if not line or line.startswith("#"): continue
            k,_,v = line.partition("=")
            env[k] = v
    return env


class SftpPool:
    """One paramiko Transport + SFTPClient per thread, lazily created."""
    def __init__(self, host, port, user, password):
        self.host, self.port, self.user, self.password = host, port, user, password
        self._local = threading.local()

    def get(self):
        c = getattr(self._local, "sftp", None)
        if c is not None:
            return c
        t = paramiko.Transport((self.host, self.port))
        t.set_keepalive(30)
        try:
            t.connect(username=self.user, password=self.password)
        except paramiko.ssh_exception.BadAuthenticationType:
            t.auth_interactive(self.user, lambda a,b,p: [self.password for _ in p])
        self._local.transport = t
        self._local.sftp = paramiko.SFTPClient.from_transport(t)
        return self._local.sftp


# ------------------------------------------------------------------
def walk_remote(sftp, root, out):
    """Recursive walk; appends (remote_path, size) tuples to out."""
    try:
        entries = sftp.listdir_attr(root)
    except Exception as e:
        print(f"[walk-err] {root}: {e}", flush=True)
        return
    for a in entries:
        p = root.rstrip("/") + "/" + a.filename
        if stat.S_ISDIR(a.st_mode or 0):
            walk_remote(sftp, p, out)
        else:
            out.append((p, a.st_size or 0))


def download_one(pool, remote, size, dest):
    """Download one file atomically. Returns (status, bytes)."""
    local = os.path.join(dest, remote.lstrip("/"))
    if os.path.exists(local) and os.path.getsize(local) == size:
        return "skip", size
    os.makedirs(os.path.dirname(local), exist_ok=True)
    tmp = local + ".tmp"
    sftp = pool.get()
    sftp.get(remote, tmp)
    if os.path.getsize(tmp) != size:
        os.remove(tmp)
        return "size_mismatch", 0
    os.replace(tmp, local)
    return "ok", size


# ------------------------------------------------------------------
def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--dest", required=True)
    ap.add_argument("--threads", type=int, default=8)
    ap.add_argument("--log", required=True)
    ap.add_argument("--env", default="/root/.kravet-sftp.env")
    ap.add_argument("--subdirs", nargs="*", default=["BF","KF","LJ"])
    args = ap.parse_args()

    env = load_env(args.env)
    pool = SftpPool(env["KRAVET_SFTP_HOST"], int(env["KRAVET_SFTP_PORT"]),
                    env["KRAVET_SFTP_USER"], env["KRAVET_SFTP_PASS"])

    # 1) Walk remote tree (single-threaded with 1 session)
    print(f"[walk] listing {args.subdirs} ...", flush=True)
    t0 = time.time()
    files = []
    sftp = pool.get()
    for sub in args.subdirs:
        walk_remote(sftp, "/" + sub, files)
    print(f"[walk] {len(files):,} files  "
          f"{sum(s for _,s in files)/1024/1024/1024:.2f} GB  "
          f"{time.time()-t0:.0f}s", flush=True)

    # 2) Download with thread pool
    os.makedirs(args.dest, exist_ok=True)
    done = 0
    skipped = 0
    failed = 0
    bytes_dl = 0
    t0 = time.time()
    lock = threading.Lock()
    log_f = open(args.log, "a", buffering=1)

    def worker(remote, size):
        try:
            status, nb = download_one(pool, remote, size, args.dest)
            return (remote, status, nb, None)
        except Exception as e:
            return (remote, "err", 0, f"{type(e).__name__}: {e}")

    with ThreadPoolExecutor(max_workers=args.threads) as ex:
        futs = [ex.submit(worker, r, s) for r,s in files]
        for f in as_completed(futs):
            remote, status, nb, err = f.result()
            with lock:
                if status == "ok":        done += 1; bytes_dl += nb
                elif status == "skip":    skipped += 1
                else:                     failed += 1
                log_f.write(json.dumps({
                    "t": int(time.time()),
                    "remote": remote,
                    "status": status,
                    "bytes": nb,
                    "error": err,
                }) + "\n")
                total_done = done + skipped + failed
                if total_done % 500 == 0:
                    el = time.time() - t0
                    rate_mb = bytes_dl / max(1, el) / 1024 / 1024
                    print(f"[{total_done:,}/{len(files):,}] "
                          f"ok={done:,} skip={skipped:,} fail={failed:,}  "
                          f"{bytes_dl/1024/1024/1024:.2f}GB  "
                          f"{rate_mb:.1f} MB/s  "
                          f"{el:.0f}s", flush=True)

    el = time.time() - t0
    rate_mb = bytes_dl / max(1, el) / 1024 / 1024
    print(f"\n[DONE] ok={done:,} skip={skipped:,} fail={failed:,}  "
          f"{bytes_dl/1024/1024/1024:.2f} GB in {el:.0f}s  "
          f"avg {rate_mb:.1f} MB/s", flush=True)
    log_f.close()


if __name__ == "__main__":
    try:
        main()
    except Exception:
        traceback.print_exc()
        sys.exit(1)