← 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)