codekingpro/portable-devtools
114k
1"""2Implementation of a custom transfer agent for the transfer type "multipart" for3git-lfs.4 5Inspired by:6github.com/cbartz/git-lfs-swift-transfer-agent/blob/master/git_lfs_swift_transfer.py7 8Spec is: github.com/git-lfs/git-lfs/blob/master/docs/custom-transfers.md9 10 11To launch debugger while developing:12 13``` [lfs "customtransfer.multipart"]14path = /path/to/huggingface_hub/.venv/bin/python args = -m debugpy --listen 567815--wait-for-client16/path/to/huggingface_hub/src/huggingface_hub/commands/huggingface_cli.py17lfs-multipart-upload ```"""18 19import json20import os21import subprocess22import sys23from typing import Annotated24 25import typer26 27from huggingface_hub.errors import CLIError28from huggingface_hub.lfs import LFS_MULTIPART_UPLOAD_COMMAND29 30from ..utils import get_session, hf_raise_for_status, logging31from ..utils._lfs import SliceFileObj32 33 34logger = logging.get_logger(__name__)35 36 37def lfs_enable_largefiles(38 path: Annotated[39 str,40 typer.Argument(41 help="Local path to repository you want to configure.",42 ),43 ],44) -> None:45 """46 Configure your repository to enable upload of files > 5GB.47 48 This command sets up git-lfs to use the custom multipart transfer agent49 which enables efficient uploading of large files in chunks.50 """51 local_path = os.path.abspath(path)52 if not os.path.isdir(local_path):53 raise CLIError("This does not look like a valid git repo.")54 subprocess.run(55 "git config lfs.customtransfer.multipart.path hf".split(),56 check=True,57 cwd=local_path,58 )59 subprocess.run(60 f"git config lfs.customtransfer.multipart.args {LFS_MULTIPART_UPLOAD_COMMAND}".split(),61 check=True,62 cwd=local_path,63 )64 print("Local repo set up for largefiles")65 66 67def write_msg(msg: dict):68 """Write out the message in Line delimited JSON."""69 msg_str = json.dumps(msg) + "\n"70 sys.stdout.write(msg_str)71 sys.stdout.flush()72 73 74def read_msg() -> dict | None:75 """Read Line delimited JSON from stdin."""76 msg = json.loads(sys.stdin.readline().strip())77 78 if "terminate" in (msg.get("type"), msg.get("event")):79 # terminate message received80 return None81 82 if msg.get("event") not in ("download", "upload"):83 logger.critical("Received unexpected message")84 sys.exit(1)85 86 return msg87 88 89def lfs_multipart_upload() -> None:90 """Internal git-lfs custom transfer agent for multipart uploads.91 92 This function implements the custom transfer protocol for git-lfs multipart uploads.93 Handles chunked uploads of large files to Hugging Face Hub.94 """95 # Immediately after invoking a custom transfer process, git-lfs96 # sends initiation data to the process over stdin.97 # This tells the process useful information about the configuration.98 init_msg = json.loads(sys.stdin.readline().strip())99 if not (init_msg.get("event") == "init" and init_msg.get("operation") == "upload"):100 write_msg({"error": {"code": 32, "message": "Wrong lfs init operation"}})101 sys.exit(1)102 103 # The transfer process should use the information it needs from the104 # initiation structure, and also perform any one-off setup tasks it105 # needs to do. It should then respond on stdout with a simple empty106 # confirmation structure, as follows:107 write_msg({})108 109 # After the initiation exchange, git-lfs will send any number of110 # transfer requests to the stdin of the transfer process, in a serial sequence.111 while True:112 msg = read_msg()113 if msg is None:114 # When all transfers have been processed, git-lfs will send115 # a terminate event to the stdin of the transfer process.116 # On receiving this message the transfer process should117 # clean up and terminate. No response is expected.118 sys.exit(0)119 120 oid = msg["oid"]121 filepath = msg["path"]122 completion_url = msg["action"]["href"]123 header = msg["action"]["header"]124 chunk_size = int(header.pop("chunk_size"))125 presigned_urls: list[str] = list(header.values())126 127 # Send a "started" progress event to allow other workers to start.128 # Otherwise they're delayed until first "progress" event is reported,129 # i.e. after the first 5GB by default (!)130 write_msg(131 {132 "event": "progress",133 "oid": oid,134 "bytesSoFar": 1,135 "bytesSinceLast": 0,136 }137 )138 139 parts = []140 with open(filepath, "rb") as file:141 for i, presigned_url in enumerate(presigned_urls):142 with SliceFileObj(143 file,144 seek_from=i * chunk_size,145 read_limit=chunk_size,146 ) as data:147 r = get_session().put(presigned_url, data=data)148 hf_raise_for_status(r)149 parts.append(150 {151 "etag": r.headers.get("etag"),152 "partNumber": i + 1,153 }154 )155 # In order to support progress reporting while data is uploading / downloading,156 # the transfer process should post messages to stdout157 write_msg(158 {159 "event": "progress",160 "oid": oid,161 "bytesSoFar": (i + 1) * chunk_size,162 "bytesSinceLast": chunk_size,163 }164 )165 166 r = get_session().post(167 completion_url,168 json={169 "oid": oid,170 "parts": parts,171 },172 )173 hf_raise_for_status(r)174 175 write_msg({"event": "complete", "oid": oid})176 