#!/usr/bin/env python3 """UpCloud: resumable, verified uploads into Jamie's Nextcloud.""" from __future__ import annotations import base64 import getpass import hashlib import json import os import posixpath import re import shlex import shutil import subprocess import sys import time from pathlib import Path VERSION = "0.1.0" CONFIG_DIR = Path.home() / ".config" / "upcloud" CONFIG_FILE = CONFIG_DIR / "config.json" DEFAULT_CONFIG = { "host": "40.160.53.29", "user": "upcloud", "ssh_key": "~/.config/upcloud/id_ed25519", "nextcloud_user": "jamiexiaojianxuan", "editor": "nano", } class UpCloudError(RuntimeError): pass def atomic_write(path: Path, data: str, mode: int = 0o600) -> None: path.parent.mkdir(parents=True, exist_ok=True) tmp = path.with_suffix(path.suffix + ".tmp") tmp.write_text(data, encoding="utf-8") os.chmod(tmp, mode) os.replace(tmp, path) def ensure_config(open_editor: bool = False) -> dict: if not CONFIG_FILE.exists(): atomic_write(CONFIG_FILE, json.dumps(DEFAULT_CONFIG, indent=2) + "\n") print(f"Created {CONFIG_FILE}") open_editor = True if open_editor: try: configured_editor = json.loads(CONFIG_FILE.read_text()).get("editor", DEFAULT_CONFIG["editor"]) except (OSError, json.JSONDecodeError): configured_editor = DEFAULT_CONFIG["editor"] editor = os.environ.get("EDITOR") or configured_editor subprocess.run(shlex.split(editor) + [str(CONFIG_FILE)], check=False) try: config = json.loads(CONFIG_FILE.read_text(encoding="utf-8")) except (OSError, json.JSONDecodeError) as exc: raise UpCloudError(f"Invalid config {CONFIG_FILE}: {exc}") from exc for key in ("host", "user", "ssh_key", "nextcloud_user"): if not config.get(key): raise UpCloudError(f"Missing config value: {key}") config["ssh_key"] = os.path.expanduser(config["ssh_key"]) if not Path(config["ssh_key"]).is_file(): raise UpCloudError(f"SSH key not found: {config['ssh_key']}") return config def ssh_base(config: dict) -> list[str]: return [ "ssh", "-T", "-o", "BatchMode=yes", "-o", "ConnectTimeout=15", "-o", "ServerAliveInterval=30", "-o", "ServerAliveCountMax=10", "-i", config["ssh_key"], f"{config['user']}@{config['host']}", ] def encode_request(payload: dict) -> str: raw = json.dumps(payload, separators=(",", ":"), ensure_ascii=False).encode() return base64.urlsafe_b64encode(raw).decode().rstrip("=") def rpc(config: dict, action: str, **kwargs) -> dict: request = {"action": action, "nextcloud_user": config["nextcloud_user"], **kwargs} cmd = ssh_base(config) + ["upcloud-rpc", encode_request(request)] proc = subprocess.run(cmd, capture_output=True, text=True, stdin=subprocess.DEVNULL) output = proc.stdout.strip() if proc.returncode != 0: detail = proc.stderr.strip() or output or f"ssh exited {proc.returncode}" raise UpCloudError(detail) try: response = json.loads(output) except json.JSONDecodeError as exc: raise UpCloudError(f"Invalid server response: {output[:300]}") from exc if not response.get("ok"): raise UpCloudError(response.get("error", "Unknown server error")) return response def normalize_remote(cwd: str, target: str) -> str: target = target.strip() if not target or target == ".": return cwd if target in ("~", "/"): return "" if target.startswith("~/"): target = target[2:] base = "" elif target.startswith("/"): target = target[1:] base = "" else: base = cwd result = posixpath.normpath(posixpath.join(base, target)) if result in (".", "/"): return "" if result == ".." or result.startswith("../"): raise UpCloudError("Cannot leave the Nextcloud root") return result.strip("/") def human_bytes(value: int) -> str: units = ["B", "KiB", "MiB", "GiB", "TiB"] number = float(value) for unit in units: if number < 1024 or unit == units[-1]: return f"{number:.1f} {unit}" if unit != "B" else f"{int(number)} B" number /= 1024 return f"{value} B" def sha256_file(path: Path) -> str: size = path.stat().st_size done = 0 digest = hashlib.sha256() started = time.monotonic() with path.open("rb") as handle: while True: chunk = handle.read(8 * 1024 * 1024) if not chunk: break digest.update(chunk) done += len(chunk) if size: percent = done * 100 / size elapsed = max(time.monotonic() - started, 0.001) rate = done / elapsed print( f"\rHashing [{'#' * int(percent / 4):<25}] {percent:6.2f}% " f"{human_bytes(rate)}/s", end="", flush=True, ) print() return digest.hexdigest() PROGRESS_RE = re.compile(r"\s*([\d,]+)\s+(\d+)%\s+([^\s]+/s)") def run_rsync(config: dict, source: Path, stage_rel: str) -> float: rsync = shutil.which("rsync") if not rsync: raise UpCloudError("rsync is not installed (macOS: brew install rsync)") version = subprocess.run([rsync, "--version"], capture_output=True, text=True) parsed = re.search(r"version\s+(\d+)\.(\d+)", version.stdout) if not parsed or tuple(map(int, parsed.groups())) < (3, 1): raise UpCloudError("rsync 3.1+ is required (macOS: brew install rsync)") ssh_transport = " ".join( shlex.quote(part) for part in [ "ssh", "-T", "-o", "BatchMode=yes", "-o", "ConnectTimeout=15", "-o", "ServerAliveInterval=30", "-o", "ServerAliveCountMax=10", "-i", config["ssh_key"], ] ) cmd = [ rsync, "-rt", "--partial", "--append-verify", "--info=progress2", "--no-compress", "-e", ssh_transport, str(source), f"{config['user']}@{config['host']}:{stage_rel}", ] started = time.monotonic() proc = subprocess.Popen( cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, errors="replace", bufsize=1, ) last = "" assert proc.stdout is not None try: while True: char = proc.stdout.read(1) if not char: break if char in "\r\n": if last: match = PROGRESS_RE.search(last) if match: percent = int(match.group(2)) bar = "#" * int(percent / 4) print(f"\rUpload [{bar:<25}] {percent:3d}% {match.group(3):>12}", end="", flush=True) elif "error" in last.lower() or "failed" in last.lower(): print(f"\n{last}") last = "" else: last += char except BaseException: proc.terminate() try: proc.wait(timeout=5) except subprocess.TimeoutExpired: proc.kill() proc.wait() raise return_code = proc.wait() print() if return_code != 0: raise UpCloudError(f"rsync failed with exit code {return_code}; rerun upload to resume") return time.monotonic() - started def list_remote(config: dict, cwd: str, target: str, show_all: bool) -> None: path = normalize_remote(cwd, target) result = rpc(config, "list", path=path, show_all=show_all) entries = result.get("entries", []) if not entries: print("(empty)") return for entry in entries: marker = "/" if entry["type"] == "dir" else "" size = "-" if entry["type"] == "dir" else human_bytes(entry.get("size", 0)) print(f"{entry['name']}{marker:<2} {size:>10}") def upload(config: dict, cwd: str, args: list[str]) -> None: if not args: raise UpCloudError('Usage: upload "local_path" [remote_name]') source = Path(os.path.expanduser(args[0])).absolute() if not source.is_file() or source.is_symlink(): raise UpCloudError("upload currently accepts one regular file at a time") remote_name = args[1] if len(args) > 1 else source.name if "/" in remote_name or remote_name in ("", ".", ".."): raise UpCloudError("remote_name must be a single file name") started_total = time.monotonic() size = source.stat().st_size print(f"File: {source}") print(f"Size: {human_bytes(size)}") local_hash = sha256_file(source) source_signature = (source.stat().st_size, source.stat().st_mtime_ns) print(f"Local SHA-256: {local_hash}") overwrite = False check = rpc(config, "exists", path=normalize_remote(cwd, remote_name)) if check.get("exists"): try: answer = input(f"{remote_name} already exists. Overwrite? [y/N] ").strip().lower() except EOFError: print("Cancelled: overwrite requires explicit confirmation.") return if answer not in ("y", "yes"): print("Cancelled.") return overwrite = True prepared = rpc( config, "prepare", path=cwd, filename=remote_name, size=size, sha256=local_hash, overwrite=overwrite, ) if prepared.get("resumed"): print("Resuming existing staged upload.") network_seconds = run_rsync(config, source, prepared["stage_rel"]) if source_signature != (source.stat().st_size, source.stat().st_mtime_ns): raise UpCloudError("Local file changed during upload; refusing to import") transferred = max(0, size - prepared.get("existing_bytes", 0)) average_mbps = (transferred * 8 / 1_000_000) / max(network_seconds, 0.001) print("Verifying staged file and importing through Nextcloud…") committed = rpc(config, "commit", job_id=prepared["job_id"]) total_seconds = time.monotonic() - started_total verified = bool(committed.get("verified")) print("\nResult") print(f" verified: {str(verified).lower()}") print(f" local SHA-256: {local_hash}") print(f" stage SHA-256: {committed.get('stage_sha256', '-')}") print(f" final SHA-256: {committed.get('final_sha256', '-')}") print(f" destination: /{committed.get('path', normalize_remote(cwd, remote_name))}") print(f" upload time: {network_seconds:.1f}s") print(f" average: {average_mbps:.2f} Mbps ({average_mbps / 8:.2f} MB/s)") print(f" new payload: {human_bytes(transferred)}") print(f" total time: {total_seconds:.1f}s") if not verified: raise UpCloudError("Final SHA-256 mismatch; staged data was retained") def show_stats(config: dict) -> None: started = time.monotonic() stats = rpc(config, "stats") rtt_ms = (time.monotonic() - started) * 1000 try: ping = subprocess.run(["ping", "-c", "3", config["host"]], capture_output=True, text=True, timeout=10) summary = next((line for line in ping.stdout.splitlines() if "min/avg/max" in line), "ICMP unavailable") except subprocess.TimeoutExpired: summary = "ICMP timeout" print(f"Ping: {summary}") print(f"SSH/API request: {rtt_ms:.1f} ms (includes server sampling)") print(f"CPU: {stats['cpu_percent']:.1f}%") print(f"Load (1/5/15m): {stats['load']}") print(f"RAM: {human_bytes(stats['ram_used'])} / {human_bytes(stats['ram_total'])} ({stats['ram_percent']:.1f}%)") print(f"/data: {human_bytes(stats['disk_used'])} / {human_bytes(stats['disk_total'])} ({stats['disk_percent']:.1f}%)") print(f"Current network: ↓ {stats['rx_mbps']:.2f} Mbps ↑ {stats['tx_mbps']:.2f} Mbps") print(f"Staging: {human_bytes(stats['staging_bytes'])} / {human_bytes(stats['staging_limit'])}") print(f"Uptime: {stats['uptime_h']:.1f} h") HELP = """Commands: ls [path] list folders/files ls -a [path] include hidden entries cd change Nextcloud folder pwd show current Nextcloud folder upload "local_path" [name] resumable upload + SHA-256 verification stats server/network statistics settings edit ~/.config/upcloud/config.json with nano help show this help exit quit """ def execute(config: dict, cwd: str, argv: list[str]) -> tuple[str, bool]: if not argv: return cwd, True command, *args = argv if command in ("exit", "quit"): return cwd, False if command == "help": print(HELP) elif command == "pwd": print("/" + cwd if cwd else "/") elif command == "ls": show_all = "-a" in args targets = [arg for arg in args if arg != "-a"] list_remote(config, cwd, targets[0] if targets else ".", show_all) elif command == "cd": destination = normalize_remote(cwd, args[0] if args else "/") result = rpc(config, "is_dir", path=destination) if not result.get("is_dir"): raise UpCloudError("Not a directory") cwd = destination elif command == "upload": upload(config, cwd, args) elif command == "stats": show_stats(config) elif command in ("setting", "settings", "config"): updated = ensure_config(open_editor=True) config.clear() config.update(updated) else: raise UpCloudError(f"Unknown command: {command}. Type help.") return cwd, True def interactive(config: dict) -> int: cwd = "" try: import readline import glob commands = ["ls", "cd", "pwd", "upload", "stats", "setting", "settings", "help", "exit"] def complete(text, state): line = readline.get_line_buffer() if " " not in line: matches = [command + " " for command in commands if command.startswith(text)] elif line.lstrip().startswith("upload "): matches = glob.glob(os.path.expanduser(text) + "*") matches = [shlex.quote(item + "/" if os.path.isdir(item) else item) for item in matches] else: if state == 0: try: parent = posixpath.dirname(text) folder = normalize_remote(cwd, parent or ".") entries = rpc(config, "list", path=folder, show_all=True)["entries"] complete.matches = [ shlex.quote((parent + "/" if parent else "") + item["name"] + ("/" if item["type"] == "dir" else "")) for item in entries if item["name"].startswith(posixpath.basename(text)) ] except UpCloudError: complete.matches = [] matches = getattr(complete, "matches", []) return matches[state] if state < len(matches) else None readline.set_completer(complete) readline.set_completer_delims(" \t\n") readline.parse_and_bind("bind ^I rl_complete" if "libedit" in (readline.__doc__ or "") else "tab: complete") except ImportError: pass print(f"UpCloud {VERSION} — verified Nextcloud uploads") print("Type help for commands.") while True: prompt = f"upcloud:/{cwd}> " if cwd else "upcloud:/> " try: line = input(prompt) except (EOFError, KeyboardInterrupt): print() return 0 try: argv = shlex.split(line) cwd, keep_running = execute(config, cwd, argv) if not keep_running: return 0 except KeyboardInterrupt: print("\nInterrupted. Re-run the same upload command to resume.") except (UpCloudError, ValueError) as exc: print(f"error: {exc}", file=sys.stderr) def main() -> int: try: if len(sys.argv) > 1 and sys.argv[1] in ("--version", "-V"): print(VERSION) return 0 config = ensure_config(open_editor=False) if len(sys.argv) == 1: return interactive(config) _, _ = execute(config, "", sys.argv[1:]) return 0 except UpCloudError as exc: print(f"error: {exc}", file=sys.stderr) return 1 except (OSError, subprocess.SubprocessError) as exc: print(f"error: {exc}", file=sys.stderr) return 1 except KeyboardInterrupt: print("\nInterrupted. Re-run the same upload command to resume.", file=sys.stderr) return 130 if __name__ == "__main__": raise SystemExit(main())