#!/usr/bin/env python3 """ migrate_s3_to_pv.py — перелив архивов функций из S3 → storagesvc (PV-backend) Запуск: kubectl exec -n fission deploy/storagesvc -- python3 /tmp/migrate.py ИЛИ: запустить как Job в кластере Что делает: 1. Читает все Package CRD во всех namespace 2. Для каждого пакета с URL s3://... скачивает архив из S3 3. Загружает в storagesvc HTTP API → получает новый URL 4. Обновляет Package CRD: .spec.deployment.url = новый URL """ import os import sys import json import time import urllib.request import urllib.parse # S3 credentials (из env или хардкод для ручного запуска) S3_ENDPOINT = os.environ.get("S3_ENDPOINT", "https://s3.msk-1.ngcloud.ru") S3_BUCKET = os.environ.get("S3_BUCKET", "sless-functions") S3_ACCESS_KEY = os.environ.get("S3_ACCESS_KEY", "0GLQRD38H4I6RBDB0EWJ") S3_SECRET_KEY = os.environ.get("S3_SECRET_KEY", "eTFibiHmBd96IApj9PYsboTR6OBoD7osxoarHykw") S3_REGION = os.environ.get("S3_REGION", "msk-1") STORAGESVC_URL = os.environ.get("STORAGESVC_URL", "http://storagesvc.fission.svc.cluster.local") import subprocess def kubectl(args, input_data=None): cmd = ["kubectl"] + args r = subprocess.run(cmd, capture_output=True, text=True, input=input_data) if r.returncode != 0: raise RuntimeError(f"kubectl {' '.join(args)} failed: {r.stderr}") return r.stdout def get_all_packages(): out = kubectl(["get", "packages", "--all-namespaces", "-o", "json"]) return json.loads(out)["items"] def is_s3_url(url): return url and url.startswith("s3://") def s3_key_from_url(url): # s3://sless-functions/fission/abc123 → fission/abc123 path = url[len(f"s3://{S3_BUCKET}/"):] return path def download_from_s3(s3_key): import boto3 s3 = boto3.client( "s3", endpoint_url=S3_ENDPOINT, aws_access_key_id=S3_ACCESS_KEY, aws_secret_access_key=S3_SECRET_KEY, region_name=S3_REGION, ) obj = s3.get_object(Bucket=S3_BUCKET, Key=s3_key) return obj["Body"].read() def upload_to_storagesvc(data): import urllib.request, uuid boundary = uuid.uuid4().hex body = ( f"--{boundary}\r\n" f'Content-Disposition: form-data; name="uploadfile"; filename="archive.zip"\r\n' f"Content-Type: application/octet-stream\r\n\r\n" ).encode() + data + f"\r\n--{boundary}--\r\n".encode() req = urllib.request.Request( f"{STORAGESVC_URL}/v1/archive", data=body, method="POST", headers={"Content-Type": f"multipart/form-data; boundary={boundary}"}, ) with urllib.request.urlopen(req, timeout=60) as resp: result = json.loads(resp.read()) return result["id"] # storagesvc возвращает {"id": "http://storagesvc.../v1/archive?id=..."} def update_package_url(ns, name, new_url): patch = json.dumps({"spec": {"deployment": {"url": new_url, "type": "url"}}}) kubectl(["patch", "package", name, "-n", ns, "--type=merge", f"--patch={patch}"]) def main(): try: import boto3 except ImportError: print("ERROR: boto3 не установлен. Запусти: pip install boto3") sys.exit(1) packages = get_all_packages() total = len(packages) migrated = 0 skipped = 0 errors = 0 print(f"Всего пакетов: {total}") print() for pkg in packages: ns = pkg["metadata"]["namespace"] name = pkg["metadata"]["name"] url = pkg.get("spec", {}).get("deployment", {}).get("url", "") if not is_s3_url(url): print(f" SKIP {ns}/{name} — url: {url or '(empty)'}") skipped += 1 continue s3_key = s3_key_from_url(url) print(f" MIG {ns}/{name} s3://{S3_BUCKET}/{s3_key}", end=" ", flush=True) try: data = download_from_s3(s3_key) new_id = upload_to_storagesvc(data) new_url = f"{STORAGESVC_URL}/v1/archive?id={new_id}" if not new_id.startswith("http") else new_id update_package_url(ns, name, new_url) print(f"→ {new_url}") migrated += 1 except Exception as e: print(f" ERROR: {e}") errors += 1 print() print(f"Готово: перенесено={migrated} пропущено={skipped} ошибок={errors}") if __name__ == "__main__": main()