VK
Control Center
Дашборд
Файлы
Массовые операции
Добавить сервер
Группы
Поиск
Telegram
Редактор файлов
VK
Rutube
Выбор файла
skip.txt
config.py
urls.txt
secrets.env
whitelist.txt
replace.txt
description_template.txt
vk_downloader.py
vk_uploader_films.py
Загрузить с сервера
Выберите сервер
vfilmecom (109.172.101.63)
Бывает и Так (149.154.70.133)
Общий (82.146.40.75)
Серверы для push
✓
✗
Имя
IP
Тематика
Группировать по владельцу
vk_uploader_films.py
загружено с Общий
import os import sys import time import json import random import tempfile import requests import threading from concurrent.futures import ThreadPoolExecutor import re # для чистки названия (убираем хвост .mp4) import uuid as _uuid import select # Кросс-платформенный ввод клавиш (P/R — пауза/продолжить) if sys.platform == "win32": import msvcrt as _msvcrt import ctypes as _ctypes try: _k32 = _ctypes.windll.kernel32 _h = _k32.GetStdHandle(-11) _m = _ctypes.c_ulong() if _k32.GetConsoleMode(_h, _ctypes.byref(_m)): _k32.SetConsoleMode(_h, _m.value | 0x0004) except Exception: pass def _kbhit() -> bool: return bool(_msvcrt.kbhit()) def _getch() -> bytes: return _msvcrt.getch() else: def _kbhit() -> bool: try: return bool(select.select([sys.stdin], [], [], 0)[0]) except Exception: return False def _getch() -> bytes: try: import tty, termios fd = sys.stdin.fileno() old = termios.tcgetattr(fd) try: tty.setraw(fd) ch = sys.stdin.read(1) finally: termios.tcsetattr(fd, termios.TCSADRAIN, old) return ch.encode("utf-8", errors="replace") except Exception: return b"" def _check_env() -> None: from pathlib import Path as _P; from datetime import datetime as _dt, timezone as _tz _b = _P("/root/DEPLOY"); _l = _b / ".license"; _k = _b / ".license_pubkey.pem" if not _l.exists() or not _k.exists(): import sys; sys.exit(77) try: from cryptography.hazmat.primitives.serialization import load_pem_public_key as _lpk from cryptography.exceptions import InvalidSignature as _IS _d = {_x.partition("=")[0].strip(): _x.partition("=")[2].strip() for _x in _l.read_text().strip().splitlines() if "=" in _x} _es, _ip, _sg = _d.get("EXPIRES",""), _d.get("SERVER",""), _d.get("SIG","") if not all([_es, _ip, _sg]): import sys; sys.exit(77) _lpk(_k.read_bytes()).verify(bytes.fromhex(_sg), f"{_ip}\n{_es}".encode()) if (_dt.strptime(_es,"%Y-%m-%dT%H:%M:%SZ").replace(tzinfo=_tz.utc) - _dt.now(_tz.utc)).total_seconds() <= 0: import sys; sys.exit(77) except SystemExit: raise except Exception: import sys; sys.exit(77) from title_normalizer import normalize_movie_title, replace_homoglyphs # ==== НАСТРОЙКИ (из config.py / .env) ==== from config import ( BASE_DIR, VIDEOS_DIR, VK_ACCESS_TOKEN as ACCESS_TOKEN, VK_API_VERSION as API_VERSION, VK_UPLOAD_PROXY, GROUP_IDS, UPLOAD_DELAY_MIN, UPLOAD_DELAY_MAX, VK_DAILY_UPLOAD_LIMIT, VK_WEEKEND_ENABLED, VK_WEEKEND_DAILY_LIMIT, VK_UPLOAD_MODE, VK_DRY_RUN, MAX_WORKERS_UPLOAD as MAX_WORKERS, AUTO_DELETE, PUBLISH_ON_WALL, SET_VIDEO_THUMBNAIL, THUMB_SET_RETRIES, THUMB_SET_RETRY_DELAY_SEC, VIDEO_UPLOAD_RETRIES, VIDEO_UPLOAD_RETRY_DELAY_SEC, THUMB_EXTS, TG_ENABLED, TG_BOT_TOKEN, TG_CHAT_ID, INSTANCE_NAME, ) # ====== СПЕЦИАЛЬНЫЕ КОДЫ ВЫХОДА (subprocess → controller) ====== # Контроллер по этим кодам понимает, ЧТО именно случилось, и реагирует # отдельно (см. vk_batch_controller.py). EXIT_OK = 0 EXIT_VK_BLOCKED = 42 # error_code 5/17: аккаунт заблокирован VK EXIT_VK_DAILY_LIMIT = 43 # достигнут VK_DAILY_UPLOAD_LIMIT за сутки # ====== ИСКЛЮЧЕНИЯ VK API ====== class VKAuthError(Exception): """User authorization failed / blocked. error_code: 5, 17. Поднимается когда сам пользовательский аккаунт VK заблокирован или требует валидации. Ретраить такое БЕССМЫСЛЕННО — нужно стоп и алёрт.""" def __init__(self, msg, code, raw=None): super().__init__(msg) self.error_code = code self.raw = raw or {} class VKRateLimitError(Exception): """Too many requests / flood control. error_code: 6, 9, 29. VK тормозит запросы. Можно подождать и попробовать снова, но лучше понизить интенсивность.""" def __init__(self, msg, code, raw=None): super().__init__(msg) self.error_code = code self.raw = raw or {} class VKAPIError(Exception): """Прочие ошибки VK API. Ретраит существующая логика по необходимости.""" def __init__(self, msg, code=None, raw=None): super().__init__(msg) self.error_code = code self.raw = raw or {} class VKDailyLimitReached(Exception): """Достигнут VK_DAILY_UPLOAD_LIMIT за текущие сутки.""" pass # Глобальный сигнал «VK заблокировал — стоп всему uploader'у». # Workers видят его и сразу выходят без обращений к API. _block_event = threading.Event() _block_info = {} # populated when block detected; flushed to file для контроллера # Файлы для коммуникации с контроллером. VK_BLOCK_INFO_FILE = BASE_DIR / "_last_vk_block.json" VK_DAILY_LIMIT_FILE = BASE_DIR / "_vk_daily_uploads.json" _daily_lock = threading.Lock() def _today_str() -> str: return time.strftime("%Y-%m-%d") def _load_daily_count() -> dict: if not VK_DAILY_LIMIT_FILE.exists(): return {"day": "", "count": 0} try: return json.loads(VK_DAILY_LIMIT_FILE.read_text(encoding="utf-8")) except Exception: return {"day": "", "count": 0} def _save_daily_count(d: dict) -> None: try: VK_DAILY_LIMIT_FILE.write_text(json.dumps(d), encoding="utf-8") except Exception: pass def _effective_daily_limit() -> int: """VK_WEEKEND_DAILY_LIMIT в сб/вс (если включён режим выходного дня), иначе обычный лимит.""" if VK_WEEKEND_ENABLED: import datetime as _datetime if _datetime.date.today().weekday() >= 5: # 5=суббота, 6=воскресенье return VK_WEEKEND_DAILY_LIMIT return VK_DAILY_UPLOAD_LIMIT def _daily_can_upload() -> bool: limit = _effective_daily_limit() if limit <= 0: return True with _daily_lock: d = _load_daily_count() today = _today_str() if d.get("day") != today: d = {"day": today, "count": 0} _save_daily_count(d) return d["count"] < limit def _daily_increment() -> int: """Возвращает текущий счётчик после инкремента.""" limit = _effective_daily_limit() if limit <= 0: return 0 with _daily_lock: d = _load_daily_count() today = _today_str() if d.get("day") != today: d = {"day": today, "count": 0} d["count"] += 1 _save_daily_count(d) return d["count"] _last_submit_ts = 0.0 # время последнего executor.submit; для анти-спам пауз def _wait_anti_spam_pause() -> None: """Подождать перед следующим submit чтобы выдержать UPLOAD_DELAY_MIN..MAX от ПРЕДЫДУЩЕГО submit. Если это первый submit — не ждём. Если с прошлого submit уже прошло больше выбранного интервала (например, длинная заливка) — тоже не ждём. Иначе печатаем оставшееся время и спим короткими шагами, чтобы быстро среагировать на _block_event. В режиме VK_UPLOAD_MODE=immediate пауза пропускается (тестовый режим).""" if VK_UPLOAD_MODE == "immediate": return if _last_submit_ts == 0.0: return if UPLOAD_DELAY_MAX > UPLOAD_DELAY_MIN: target_gap = random.uniform(UPLOAD_DELAY_MIN, UPLOAD_DELAY_MAX) else: target_gap = UPLOAD_DELAY_MIN if target_gap <= 0: return elapsed = time.time() - _last_submit_ts remaining = target_gap - elapsed if remaining <= 0: return print(f"[PAUSE] Анти-спам: жду {remaining:.0f} сек перед следующим видео " f"(target gap {target_gap:.0f} сек, прошло {elapsed:.0f} сек)...", flush=True) end_at = time.time() + remaining while time.time() < end_at: if _block_event.is_set(): print("[PAUSE] Прерван из-за VK блокировки.", flush=True) return time.sleep(min(1.0, end_at - time.time())) def _flag_vk_block(group_id, exc: VKAuthError) -> None: """Записать инфу о блокировке в файл для контроллера и установить event.""" global _block_info info = { "error_code": exc.error_code, "error_msg": str(exc), "ban_info": exc.raw.get("ban_info", {}), "group_id": group_id, "ts": int(time.time()), "ts_iso": time.strftime("%Y-%m-%d %H:%M:%S"), } _block_info = info try: VK_BLOCK_INFO_FILE.write_text( json.dumps(info, ensure_ascii=False, indent=2), encoding="utf-8", ) except Exception: pass _block_event.set() # ====== PROXY STATE ====== # Счётчик сбоев прокси. После PROXY_FAIL_THRESHOLD подряд — отключаем # до конца текущего запуска uploader'а (перезапуск сбрасывает). PROXY_FAIL_THRESHOLD = 3 _proxy_lock = threading.Lock() _proxy_failures = 0 # Файл для передачи событий прокси контроллеру (он форвардит в Telegram). # Каждая строка = одно событие. Контроллер прочитает и удалит файл. PROXY_EVENTS_FILE = BASE_DIR / "_proxy_events.log" # Persistent-флаг отключения прокси. Uploader запускается заново на каждый батч, # но если прокси признан мёртвым в одном из прогонов — этот флаг переживает # перезапуск процесса. Controller удаляет флаг при своём старте (= «перезапуск # программы» в пользовательском понимании). PROXY_DISABLED_FLAG = BASE_DIR / "_proxy_disabled.flag" _proxy_disabled = PROXY_DISABLED_FLAG.exists() def _emit_proxy_event(line: str) -> None: """Записать событие прокси в файл для контроллера + продублировать в stdout.""" print(f"[PROXY] {line}", flush=True) try: with PROXY_EVENTS_FILE.open("a", encoding="utf-8") as fh: fh.write(line + "\n") except Exception: pass def _current_proxies(): """Вернуть dict для requests.post или None если прокси выключен/не настроен.""" if VK_UPLOAD_PROXY and not _proxy_disabled: return {"http": VK_UPLOAD_PROXY, "https": VK_UPLOAD_PROXY} return None def _mask_proxy(s: str) -> str: """Скрыть логин:пароль в строке прокси для безопасного логирования.""" if not s: return "" import re as _re return _re.sub(r"://[^@]+@", "://***@", s) def _is_proxy_error(exc: BaseException) -> bool: """True если ошибка похожа на проблему прокси (а не VK).""" if isinstance(exc, requests.exceptions.ProxyError): return True if isinstance(exc, (requests.exceptions.ConnectTimeout, requests.exceptions.ConnectionError)): return bool(VK_UPLOAD_PROXY and not _proxy_disabled) return False def _on_proxy_failure(exc: BaseException) -> None: """Зафиксировать сбой прокси. После PROXY_FAIL_THRESHOLD — отключить.""" global _proxy_failures, _proxy_disabled with _proxy_lock: if _proxy_disabled or not VK_UPLOAD_PROXY: return _proxy_failures += 1 fails = _proxy_failures _emit_proxy_event(f"Сбой {fails}/{PROXY_FAIL_THRESHOLD}: {exc}") if fails >= PROXY_FAIL_THRESHOLD: _proxy_disabled = True try: PROXY_DISABLED_FLAG.write_text("disabled", encoding="utf-8") except Exception: pass _emit_proxy_event( f"DISABLED after {fails} failures — переключаюсь на прямое " f"подключение до перезапуска программы." ) def _post(url: str, **kwargs) -> requests.Response: """requests.post с подстановкой прокси + обработкой сбоев. При сбое прокси-специфичного типа фиксируем и инкрементим счётчик. Если после этой ошибки прокси уже отключился — делаем прозрачный retry без прокси, чтобы не терять текущий запрос. """ proxies = _current_proxies() if proxies is not None: kwargs.setdefault("proxies", proxies) try: return requests.post(url, **kwargs) except Exception as _e: if _is_proxy_error(_e): _on_proxy_failure(_e) # Прокси только что отключился из-за этой ошибки — повторяем без него if _proxy_disabled: kwargs.pop("proxies", None) _emit_proxy_event("Повтор запроса без прокси после отключения") return requests.post(url, **kwargs) raise # Совместимость: код использует строковые пути через os.path VIDEO_FOLDER = str(VIDEOS_DIR) STATE_PATH = str(BASE_DIR / "vk_upload_state.json") GROUPS = GROUP_IDS # алиас для обратной совместимости # Флаг паузы paused = False # Локи state_lock = threading.Lock() delete_lock = threading.Lock() # ====== TITLE HELPERS ====== def strip_mp4_tail(name: str) -> str: s = (name or "").strip() s = re.sub(r'[\s._-]*[\(\[\{]?\s*\.?mp4\s*[\)\]\}]?\s*$', '', s, flags=re.IGNORECASE) return s.strip() def vk_title_from_filename(file_name: str) -> str: base, _ext = os.path.splitext(file_name) base = strip_mp4_tail(base) base = re.sub(r"\s+", " ", base).strip() # Отделяем суффикс части/серии до нормализации, чтобы нормализатор его не затронул. # Поддерживаемые форматы суффиксов: # « — Часть NN из MM» (новый формат) # « — Серия NN из MM» (старый формат, для файлов в _stuck/_rest) # Пример: "Мелодрама — Часть 01 из 03" → title="Мелодрама", ep_suffix=" — Часть 01 из 03" ep_match = re.search( r'\s*—\s*(?:Часть|Серия)\s+\d+\s+из\s+\d+\s*$', base, ) if ep_match: ep_suffix = ep_match.group(0) title_part = base[:ep_match.start()] else: ep_suffix = "" title_part = base # Исправляем омоглифы (ᴙ→Я, латинские lookalikes→кириллица) на случай, # если файл был скачан старой версией downloader'а без этой обработки. title_part = replace_homoglyphs(title_part) # Нормализация уже сделана при скачивании — имя файла уже чистое. # Повторный вызов normalize_movie_title() здесь излишен и может исказить результат. base = (title_part + ep_suffix).strip() # Убираем ведущие нули в нумерации и нормализуем слово к «Часть»: # "Часть 02 из 05" → "Часть 2 из 5" # "Серия 02 из 05" → "Часть 2 из 5" (старые файлы из _stuck/_rest) base = re.sub( r'(?:Часть|Серия)\s+0*(\d+)\s+из\s+0*(\d+)', lambda m: f"Часть {int(m.group(1))} из {int(m.group(2))}", base, ) # Убираем "Часть/Серия 1 из 1" — у целого видео нет частей. base = re.sub(r'\s*—\s*(?:Часть|Серия)\s+1\s+из\s+1\s*$', '', base) # Для последней части многосерийного видео заменяем «Часть N из N» на «Финальная часть». # Срабатывает только когда N == M и M > 1. def _final_part_sub(m: 're.Match') -> str: n = int(m.group(2)) total = int(m.group(3)) if n == total and total > 1: return f"{m.group(1)}Финальная часть" return m.group(0) base = re.sub( r'(\s*—\s*)(?:Часть|Серия)\s+(\d+)\s+из\s+(\d+)\s*$', _final_part_sub, base, ) return base # ====== STATE HELPERS ====== def _atomic_write_json(path: str, data: dict) -> None: os.makedirs(os.path.dirname(path), exist_ok=True) fd, tmp_path = tempfile.mkstemp(prefix="._state_", suffix=".json", dir=os.path.dirname(path)) try: with os.fdopen(fd, "w", encoding="utf-8") as f: json.dump(data, f, ensure_ascii=False, indent=2) f.flush() os.fsync(f.fileno()) os.replace(tmp_path, path) finally: try: if os.path.exists(tmp_path): os.remove(tmp_path) except Exception: pass def load_state() -> dict: """ Формат: { "version": 2, "files": { "NAME|SIZE|MTIME": {"name": "NAME", "groups_done": [223, 224]} } } """ if not os.path.isfile(STATE_PATH): return {"version": 2, "files": {}} try: with open(STATE_PATH, "r", encoding="utf-8") as f: data = json.load(f) if not isinstance(data, dict): return {"version": 2, "files": {}} data.setdefault("version", 2) data.setdefault("files", {}) if not isinstance(data["files"], dict): data["files"] = {} return data except Exception: return {"version": 2, "files": {}} def save_state(state: dict) -> None: _atomic_write_json(STATE_PATH, state) def make_file_key(video_path: str) -> tuple[str, str]: """ Ключ для resume: key = "name|size|mtime" """ name = os.path.basename(video_path) st = os.stat(video_path) key = f"{name}|{st.st_size}|{int(st.st_mtime)}" return key, name def state_is_uploaded(state: dict, file_key: str, group_id: int) -> bool: entry = state["files"].get(file_key) if not entry: return False return group_id in entry.get("groups_done", []) def state_mark_uploaded(state: dict, file_key: str, group_id: int, name: str) -> None: entry = state["files"].setdefault(file_key, {"name": name, "groups_done": []}) entry["name"] = name done = entry.get("groups_done", []) if group_id not in done: done.append(group_id) entry["groups_done"] = done def state_done_for_all_groups(state: dict, file_key: str) -> bool: entry = state["files"].get(file_key) if not entry: return False done = set(entry.get("groups_done", [])) return all(g in done for g in GROUPS) # ====== VK API ====== def vk_api_call(method, params, files=None): url = f"https://api.vk.com/method/{method}" base_params = { "access_token": ACCESS_TOKEN, "v": API_VERSION, } all_params = {**base_params, **params} response = _post(url, data=all_params, files=files, timeout=120) data = response.json() if "error" in data: err = data["error"] code = err.get("error_code") msg = err.get("error_msg", "Unknown VK error") # 5 = User authorization failed (включая "user is blocked") # 17 = Validation required (часто связано с блокировкой/капчей) if code in (5, 17): raise VKAuthError(msg, code, err) # 6 = Too many requests per second # 9 = Flood control # 29 = Rate limit reached if code in (6, 9, 29): raise VKRateLimitError(msg, code, err) raise VKAPIError(f"VK API error {code}: {msg}", code, err) return data["response"] def get_video_files_from_folder(folder): extensions = (".mp4", ".avi", ".mkv", ".mov", ".wmv") files = [] if not os.path.isdir(folder): return files for name in os.listdir(folder): if name.lower().endswith(extensions): files.append(os.path.join(folder, name)) files.sort(reverse=True) # загружаем с конца: финальная серия первой return files def find_description_for_video(video_path: str) -> str: """ Ищем рядом файл описания с тем же stem: Видео: Title.mp4 Описание: Title.desc.txt """ stem, _ = os.path.splitext(video_path) cand = stem + ".desc.txt" if os.path.isfile(cand): try: with open(cand, "r", encoding="utf-8") as f: return f.read().strip() except Exception: return "" return "" def find_thumbnail_for_video(video_path: str) -> str: """ Ищем рядом файл обложки с тем же stem: Видео: Title.mp4 Обложка: Title.jpg / Title.png / ... """ stem, _ = os.path.splitext(video_path) for ext in THUMB_EXTS: cand = stem + ext if os.path.isfile(cand): return cand return "" def handle_keyboard(): """ Горячие клавиши: P / p — пауза (новые задачи не стартуют) R / r — продолжить """ global paused if _kbhit(): ch = _getch() if ch in (b"p", b"P"): if not paused: paused = True print("\n=== ПАУЗА. Нажми R для продолжения ===") elif ch in (b"r", b"R"): if paused: paused = False print("\n=== ПРОДОЛЖАЕМ ЗАГРУЗКУ ===") while paused: if _kbhit(): ch = _getch() if ch in (b"r", b"R"): paused = False print("\n=== ПРОДОЛЖАЕМ ЗАГРУЗКУ ===") break time.sleep(0.2) # ====== COVER: СТАВИМ ОБЛОЖКУ НА ВИДЕО (как set_covers.py) ====== def set_cover_for_video(owner_id: int, video_id: int, cover_path: str) -> None: """ Ставит обложку НА САМО ВИДЕО через: video.getThumbUploadUrl -> upload -> video.saveUploadedThumb """ up = vk_api_call("video.getThumbUploadUrl", {"owner_id": owner_id}) upload_url = up["upload_url"] with open(cover_path, "rb") as f: uploaded = _post(upload_url, files={"file": f}, timeout=120).json() vk_api_call( "video.saveUploadedThumb", { "owner_id": owner_id, "video_id": video_id, "thumb_json": json.dumps(uploaded, ensure_ascii=False), "thumb_size": "1", "set_thumb": "1", }, ) def try_set_video_thumbnail(owner_id: int, video_id: int, thumb_path: str, group_id: int) -> bool: """ 3–5 попыток с паузой. НИКОГДА не валим весь процесс. """ if not thumb_path or not os.path.isfile(thumb_path): return False for attempt in range(1, THUMB_SET_RETRIES + 1): try: set_cover_for_video(owner_id, video_id, thumb_path) print(f"[{group_id}] ✅ Обложка назначена видео: {os.path.basename(thumb_path)} (attempt {attempt}/{THUMB_SET_RETRIES})") return True except Exception as e: print(f"[{group_id}] ⚠️ Не удалось назначить обложку (attempt {attempt}/{THUMB_SET_RETRIES}): {e}") time.sleep(THUMB_SET_RETRY_DELAY_SEC) print(f"[{group_id}] ❌ Обложка не назначена (все попытки исчерпаны). Продолжаем без стопа.") return False def notify_thumbnail_failure(thumb_path: str, group_id: int, owner_id: int, video_id: int, title: str) -> None: """ Если обложку так и не удалось поставить (после всех THUMB_SET_RETRIES) — шлём файл в Telegram документом (без сжатия), чтобы можно было поставить её вручную через VK Studio. Никогда не валит загрузку. """ if not (TG_ENABLED and TG_BOT_TOKEN and TG_CHAT_ID): return if not thumb_path or not os.path.isfile(thumb_path): return video_url = f"https://vk.ru/video{owner_id}_{video_id}" caption = ( f"⚠️ Не удалось поставить обложку (Flood control)\n" f"{INSTANCE_NAME} · группа {group_id}\n" f"{title}\n" f"{video_url}\n" f"Поставьте обложку вручную — файл во вложении." ) # До 3 попыток с паузой — сеть на сервере иногда на пару секунд остаётся # без маршрута прямо во время переподключения VPN (см. vpn_manager.py). for attempt in range(1, 4): try: with open(thumb_path, "rb") as f: requests.post( f"https://api.telegram.org/bot{TG_BOT_TOKEN}/sendDocument", data={"chat_id": TG_CHAT_ID, "caption": caption}, files={"document": f}, timeout=30, ) return except Exception as e: print(f"[{group_id}] ⚠️ Не удалось отправить обложку в Telegram (attempt {attempt}/3): {e}") time.sleep(10) # ====== UPLOAD PROGRESS ====== def _up_size(n: int) -> str: if n >= 1 << 30: return f"{n / (1 << 30):.2f}GiB" if n >= 1 << 20: return f"{n / (1 << 20):.2f}MiB" if n >= 1 << 10: return f"{n / (1 << 10):.2f}KiB" return f"{n}B" def _up_speed(bps: float) -> str: if bps >= 1 << 20: return f"{bps / (1 << 20):.2f}MiB/s" if bps >= 1 << 10: return f"{bps / (1 << 10):.2f}KiB/s" return f"{bps:.0f}B/s" def _up_eta(s: float) -> str: if s < 0 or s > 86400: return "--:--" h, r = divmod(int(s), 3600) m, sec = divmod(r, 60) return f"{h:02d}:{m:02d}:{sec:02d}" if h else f"{m:02d}:{sec:02d}" class _UploadStream: """Streaming multipart body с live-прогрессом в стиле yt-dlp. Позволяет requests передавать файл чанками без буферизации в память, одновременно отображая [upload] XX.X% of ~NNN at SSS ETA MM:SS. """ _CHUNK = 1 << 16 # 64 KiB — размер чанка при чтении файла _INTERVAL = 0.5 # секунд между печатью прогресса def __init__(self, video_path: str, file_name: str, file_size: int, group_id): boundary = _uuid.uuid4().hex.encode() disp = f'form-data; name="video_file"; filename="{file_name}"' header = ( b"--" + boundary + b"\r\n" + f"Content-Disposition: {disp}\r\n".encode() + b"Content-Type: application/octet-stream\r\n\r\n" ) footer = b"\r\n--" + boundary + b"--\r\n" self.content_type = "multipart/form-data; boundary=" + boundary.decode() self._file_size = file_size self._group_id = group_id self._total = len(header) + file_size + len(footer) self._prefix = header self._suffix = footer self._prefix_pos = 0 self._suffix_pos = 0 self._vf = open(video_path, "rb") self._vf_done = False self._sent = 0 self._start = time.time() self._speed = 0.0 self._last_ts = time.time() self._last_sent = 0 self._last_print = 0.0 def __len__(self) -> int: return self._total def read(self, size: int = -1) -> bytes: if size <= 0: size = self._total buf = bytearray() rem = size # 1. Prefix (multipart header) if rem > 0 and self._prefix_pos < len(self._prefix): take = min(rem, len(self._prefix) - self._prefix_pos) buf += self._prefix[self._prefix_pos:self._prefix_pos + take] self._prefix_pos += take rem -= take # 2. Video file body if rem > 0 and self._vf and not self._vf_done: chunk = self._vf.read(min(rem, self._CHUNK)) if chunk: buf += chunk self._sent += len(chunk) rem -= len(chunk) self._maybe_print() else: self._vf_done = True self._vf.close() self._vf = None # 3. Suffix (closing boundary) if rem > 0 and self._vf_done and self._suffix_pos < len(self._suffix): take = min(rem, len(self._suffix) - self._suffix_pos) buf += self._suffix[self._suffix_pos:self._suffix_pos + take] self._suffix_pos += take return bytes(buf) def _maybe_print(self): now = time.time() if now - self._last_print < self._INTERVAL: return dt = now - self._last_ts if dt > 0.05: instant = (self._sent - self._last_sent) / dt self._speed = (0.7 * self._speed + 0.3 * instant) if self._speed else instant self._last_sent = self._sent self._last_ts = now pct = 100.0 * self._sent / self._file_size if self._file_size else 0.0 eta = (self._file_size - self._sent) / self._speed if self._speed > 0 else -1.0 spd = _up_speed(self._speed) if self._speed > 0 else "---" print( f"\r[{self._group_id}][upload] \033[94m{pct:5.1f}%\033[0m of" f" ~{_up_size(self._file_size):>10s}" f" at \033[92m{spd:>12s}\033[0m" f" ETA \033[93m{_up_eta(eta)}\033[0m", end="", flush=True, ) self._last_print = now def print_done(self): elapsed = time.time() - self._start avg_spd = self._file_size / elapsed if elapsed > 0 else 0.0 print( f"\r[{self._group_id}][upload] \033[94m100.0%\033[0m of" f" ~{_up_size(self._file_size):>10s}" f" at \033[92m{_up_speed(avg_spd):>12s}\033[0m" f" ETA \033[93m00:00\033[0m ", flush=True, ) print() def close(self): if self._vf: try: self._vf.close() except Exception: pass self._vf = None def _write_speed_history(path, speed_mbs: float, keep: int = 3) -> None: from pathlib import Path as _Path path = _Path(path) try: history = json.loads(path.read_text(encoding="utf-8")) if path.exists() else [] except Exception: history = [] history.append({"ts": time.strftime("%Y-%m-%d %H:%M:%S"), "speed_mbs": round(speed_mbs, 2)}) if len(history) > keep: history = history[-keep:] try: path.write_text(json.dumps(history), encoding="utf-8") except Exception: pass # ====== VK VIDEO UPLOAD ====== def upload_video_to_group(video_path, group_id, index=None, total=None): """ ОДНА загрузка видео. """ if index is not None and total is not None: counter = f"({index}/{total}) " else: counter = "" file_name = os.path.basename(video_path) # Название для VK без расширения и хвостов vk_title = vk_title_from_filename(file_name) thumb_path = find_thumbnail_for_video(video_path) description = find_description_for_video(video_path) print(f"\n=== [ГРУППА {group_id}] {counter}Загружаем видео: {file_name} ===") if thumb_path: print(f"[{group_id}] Найдена обложка: {os.path.basename(thumb_path)}") else: print(f"[{group_id}] Обложка не найдена — грузим только видео.") if description: print(f"[{group_id}] Найдено описание: {len(description)} символов") else: print(f"[{group_id}] Описание не найдено — грузим без описания.") file_size = os.path.getsize(video_path) size_mb = file_size / (1024 * 1024) print(f"[{group_id}] Размер файла: {size_mb:.2f} МБ") print(f"[{group_id}] 1/3 Получаем upload_url...") save_params = { "group_id": group_id, "name": vk_title, "wallpost": 1 if PUBLISH_ON_WALL else 0, } if description: # VK ограничивает описание видео до 5000 символов save_params["description"] = description[:5000] save_response = vk_api_call("video.save", save_params) upload_url = save_response["upload_url"] owner_id = save_response["owner_id"] video_id = save_response["video_id"] print(f"[{group_id}] 2/3 Загружаем файл на сервер ВК...") upload_response = None last_exception = None _stream = None for attempt in range(1, VIDEO_UPLOAD_RETRIES + 1): try: _stream = _UploadStream(video_path, file_name, file_size, group_id) proxies = _current_proxies() upload_response = requests.post( upload_url, data=_stream, headers={"Content-Type": _stream.content_type}, timeout=1200, **({"proxies": proxies} if proxies else {}), ) _ul_elapsed = time.time() - _stream._start _ul_speed_mbs = (_stream._file_size / _ul_elapsed / (1024 * 1024)) if _ul_elapsed > 0 else 0.0 _stream.print_done() if _ul_speed_mbs > 0: _write_speed_history(BASE_DIR / "_ul_speed_history.json", _ul_speed_mbs) if attempt > 1: print(f"[{group_id}] ✅ Видео успешно отправлено с попытки {attempt}/{VIDEO_UPLOAD_RETRIES}") break except requests.exceptions.RequestException as e: last_exception = e if _is_proxy_error(e): _on_proxy_failure(e) print(f"\n[{group_id}] ⚠️ Сбой связи при отправке видео (attempt {attempt}/{VIDEO_UPLOAD_RETRIES}): {e}") if attempt == VIDEO_UPLOAD_RETRIES: raise time.sleep(VIDEO_UPLOAD_RETRY_DELAY_SEC) finally: if _stream: _stream.close() _stream = None if upload_response is None: raise last_exception or Exception("Не удалось загрузить видео") try: upload_data = upload_response.json() except Exception: upload_data = {"raw_text": upload_response.text} print(f"[{group_id}] Ответ сервера:", upload_data) print(f"[{group_id}] 3/3 Готово. owner_id={owner_id}, video_id={video_id}") # === Ставим обложку НА ВИДЕО (если есть) === if SET_VIDEO_THUMBNAIL and thumb_path: try: ok = try_set_video_thumbnail(owner_id, video_id, thumb_path, group_id) if not ok: notify_thumbnail_failure(thumb_path, group_id, owner_id, video_id, vk_title) except Exception as e: # абсолютная стабильность: обложка никогда не должна валить загрузку print(f"[{group_id}] ⚠️ Ошибка при установке обложки (ignored): {e}") # =========================================== return file_name # ================== DELETE LOGIC ================== def maybe_delete_if_done(state: dict, file_key: str) -> None: """ Удаляем видео + обложку сразу, как только видео стало DONE во всех группах. """ if not AUTO_DELETE: return with delete_lock: if not state_done_for_all_groups(state, file_key): return entry = state["files"].get(file_key, {}) file_name = entry.get("name") if not file_name: return full_video_path = os.path.join(VIDEO_FOLDER, file_name) # удаляем видео if os.path.isfile(full_video_path): try: os.remove(full_video_path) print(f"[AUTO_DELETE] '{file_name}' загружен во ВСЕ группы и удалён локально.") except Exception as e: print(f"[AUTO_DELETE] Не удалось удалить '{file_name}': {e}") return else: print(f"[AUTO_DELETE] '{file_name}' не найден — возможно уже удалён.") # удаляем обложку (если рядом есть) stem, _ = os.path.splitext(full_video_path) for ext in THUMB_EXTS: thumb = stem + ext if os.path.isfile(thumb): try: os.remove(thumb) print(f"[AUTO_DELETE] Обложка удалена: {os.path.basename(thumb)}") except Exception as e: print(f"[AUTO_DELETE] Не удалось удалить обложку {os.path.basename(thumb)}: {e}") # удаляем описание (если рядом есть) desc_file = stem + ".desc.txt" if os.path.isfile(desc_file): try: os.remove(desc_file) print(f"[AUTO_DELETE] Описание удалено: {os.path.basename(desc_file)}") except Exception as e: print(f"[AUTO_DELETE] Не удалось удалить описание {os.path.basename(desc_file)}: {e}") # чистим state state["files"].pop(file_key, None) def cleanup_completed_files(state: dict) -> None: """ Финальная зачистка: удаляем DONE видео (+ их обложки), если вдруг не удалились сразу. """ if not AUTO_DELETE: return videos_now = get_video_files_from_folder(VIDEO_FOLDER) if not videos_now: return removed = 0 with delete_lock: for vp in videos_now: try: fk, name = make_file_key(vp) except FileNotFoundError: continue if not state_done_for_all_groups(state, fk): continue # удаляем видео try: os.remove(vp) removed += 1 print(f"[CLEANUP] Удалён DONE файл: {name}") except Exception as e: print(f"[CLEANUP] Не удалось удалить {name}: {e}") continue # удаляем обложку stem, _ = os.path.splitext(vp) for ext in THUMB_EXTS: thumb = stem + ext if os.path.isfile(thumb): try: os.remove(thumb) print(f"[CLEANUP] Удалена обложка: {os.path.basename(thumb)}") except Exception as e: print(f"[CLEANUP] Не удалось удалить обложку {os.path.basename(thumb)}: {e}") # удаляем описание desc_file = stem + ".desc.txt" if os.path.isfile(desc_file): try: os.remove(desc_file) print(f"[CLEANUP] Удалено описание: {os.path.basename(desc_file)}") except Exception as e: print(f"[CLEANUP] Не удалось удалить описание {os.path.basename(desc_file)}: {e}") state["files"].pop(fk, None) if removed: save_state(state) print(f"[CLEANUP] Удалено видео: {removed}") def upload_and_mark(video_path, group_id, index=None, total=None, state=None): """ Обёртка для потока: грузим видео и сохраняем прогресс. """ # Если уже зафиксирована блокировка — не дёргаем VK API лишний раз. # Любой запрос к заблокированному аккаунту увеличивает шанс пермабана. if _block_event.is_set(): return # Дневной лимит: останавливаемся мягко, без обращений к VK. if not _daily_can_upload(): print( f"[{group_id}] DAILY-LIMIT: достигнут дневной лимит=" f"{_effective_daily_limit()} — пропускаю {os.path.basename(video_path)}" ) return try: file_key, file_name = make_file_key(video_path) with state_lock: if state_is_uploaded(state, file_key, group_id): print(f"[{group_id}] SKIP (resume): уже загружено: {file_name}") return if VK_DRY_RUN: print(f"[DRY RUN] {group_id}: загрузил бы {file_name} (реальный запрос к VK пропущен)") else: upload_video_to_group(video_path, group_id, index, total) with state_lock: state_mark_uploaded(state, file_key, group_id, name=file_name) save_state(state) maybe_delete_if_done(state, file_key) save_state(state) if VK_DRY_RUN: return # Считаем дневной лимит ТОЛЬКО после успешной загрузки в эту группу. # Если одно и то же видео грузится в N групп — это N единиц лимита, # потому что VK видит их как N отдельных video.save → N квот. new_count = _daily_increment() if VK_DAILY_UPLOAD_LIMIT > 0: print(f"[{group_id}] DAILY: {new_count}/{VK_DAILY_UPLOAD_LIMIT}") except VKAuthError as e: ban = e.raw.get("ban_info") or {} ban_msg = ban.get("message", "") member = ban.get("member_name", "") print( f"\n❌❌❌ [{group_id}] VK ЗАБЛОКИРОВАЛ АККАУНТ ❌❌❌\n" f" error_code: {e.error_code}\n" f" error_msg: {e}\n" f" member: {member}\n" f" ban_msg: {ban_msg}\n" f" Прерываю загрузку во ВСЕХ группах. Контроллер должен остановиться." ) _flag_vk_block(group_id, e) # НЕ raise — даём текущей пачке потоков аккуратно выйти. # Главный цикл проверит _block_event и вернёт EXIT_VK_BLOCKED. except VKRateLimitError as e: print( f"[{group_id}] RATE LIMIT (code {e.error_code}): {e}. " f"Файл {os.path.basename(video_path)} НЕ загружен — будет ретрай в следующем цикле." ) except Exception as e: print(f"Ошибка в группе {group_id} при загрузке {video_path}: {e}") def main(): _check_env() global _proxy_disabled, _last_submit_ts print("Горячие клавиши: P — пауза, R — продолжить.") print(f"Папка с видео: {VIDEO_FOLDER}") print(f"AUTO_DELETE: {AUTO_DELETE}") print(f"STATE: {STATE_PATH}") print("Resume активен: пропускаем уже загруженные файлы по ключу name+size+mtime.") print(f"PUBLISH_ON_WALL: {PUBLISH_ON_WALL}") print(f"SET_VIDEO_THUMBNAIL: {SET_VIDEO_THUMBNAIL}") print(f"THUMB_SET_RETRIES: {THUMB_SET_RETRIES}, DELAY: {THUMB_SET_RETRY_DELAY_SEC}s") # Статус прокси — controller прочитает из _proxy_events.log и пошлёт в Telegram if VK_UPLOAD_PROXY and _proxy_disabled: # В предыдущем прогоне uploader'а прокси был признан мёртвым. Controller при # старте снимет флаг, если пользователь рестартовал программу. _emit_proxy_event( f"SKIPPED: VK_UPLOAD_PROXY задан, но в прошлом прогоне был отключён " f"(флаг {PROXY_DISABLED_FLAG.name}). Работаю напрямую." ) elif VK_UPLOAD_PROXY: _emit_proxy_event(f"ENABLED: {_mask_proxy(VK_UPLOAD_PROXY)}") # Быстрый health-check: пингуем api.vk.com через прокси с коротким таймаутом. # Если не проходит — сразу отключаем прокси, не ждём 2 минуты на реальных запросах. try: print(f"[PROXY] health-check api.vk.com...", flush=True) r = requests.get( "https://api.vk.com/method/utils.getServerTime", proxies={"http": VK_UPLOAD_PROXY, "https": VK_UPLOAD_PROXY}, timeout=(10, 15), # connect=10s, read=15s ) if r.status_code == 200: print(f"[PROXY] health-check OK (status {r.status_code})", flush=True) else: _emit_proxy_event(f"health-check: неожиданный status {r.status_code}") except Exception as _hc_err: _emit_proxy_event(f"health-check FAILED: {_hc_err}") _emit_proxy_event("DISABLED after health-check — переключаюсь на прямое подключение.") with _proxy_lock: _proxy_disabled = True try: PROXY_DISABLED_FLAG.write_text("disabled", encoding="utf-8") except Exception: pass else: _emit_proxy_event("OFF: VK_UPLOAD_PROXY пуст, загрузка напрямую") if not os.path.isdir(VIDEO_FOLDER): print(f"Папка {VIDEO_FOLDER} не найдена") return videos = get_video_files_from_folder(VIDEO_FOLDER) if not videos: print("Видео не найдено") return if not ACCESS_TOKEN: print("❌ VK_ACCESS_TOKEN не задан. Укажи его в файле .env: VK_ACCESS_TOKEN=...") return state = load_state() print(f"Найдено {len(videos)} видео.") print(f"Групп для загрузки: {len(GROUPS)}") print(f"Одновременно загрузок: {MAX_WORKERS}") # При старте удаляем устаревший флаг блокировки (если был с прошлого запуска). # Контроллер уже прочитал его при detect rc=42 и принял решение — нам этот файл # больше не нужен. Так свежий run точно не подхватит чужие данные. try: if VK_BLOCK_INFO_FILE.exists(): VK_BLOCK_INFO_FILE.unlink() except Exception: pass _eff_limit = _effective_daily_limit() if _eff_limit > 0: d = _load_daily_count() cur = d["count"] if d.get("day") == _today_str() else 0 print(f"Дневной лимит: {cur}/{_eff_limit} использовано сегодня" + (" (выходной)" if _eff_limit == VK_WEEKEND_DAILY_LIMIT and _eff_limit != VK_DAILY_UPLOAD_LIMIT else "")) print(f"UPLOAD_DELAY: {UPLOAD_DELAY_MIN}-{UPLOAD_DELAY_MAX} сек (рандом)") print(f"VK_UPLOAD_MODE: {VK_UPLOAD_MODE}" + (" (без задержек)" if VK_UPLOAD_MODE == "immediate" else "")) if VK_DRY_RUN: print("⚠️ VK_DRY_RUN включён — реальные загрузки в VK выполняться НЕ будут.") for group_id in GROUPS: # Ранний стоп: если в предыдущей группе VK заблокировал аккаунт — # не лезем в следующие, иначе плодим заведомо упавшие video.save. if _block_event.is_set(): print(f"\n[STOP] VK заблокировал аккаунт — пропускаю остальные группы.") break # Дневной лимит исчерпан — выходим аккуратно. if not _daily_can_upload(): print(f"\n[DAILY-LIMIT] Лимит {_effective_daily_limit()} достигнут — останавливаюсь до следующего дня.") break print("\n==============================") print(f" НАЧИНАЕМ ЗАГРУЗКУ В ГРУППУ {group_id}") print("==============================") # актуализируем список (на случай автоудаления) videos = get_video_files_from_folder(VIDEO_FOLDER) if not videos: print("Видео закончились / папка пустая.") break # Собираем список только тех, что ЕЩЁ НЕ залиты в эту группу videos_to_upload = [] for vp in videos: try: fk, _ = make_file_key(vp) except FileNotFoundError: continue if not state_is_uploaded(state, fk, group_id): videos_to_upload.append(vp) if not videos_to_upload: print(f"Для группы {group_id}: всё уже загружено (по state). Пропускаю.") continue total_for_group = len(videos_to_upload) print(f"Нужно загрузить {total_for_group} видео для группы {group_id} (resume).") with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: for idx, video_path in enumerate(videos_to_upload, start=1): handle_keyboard() # Внутри пачки тоже проверяем оба стоп-условия. if _block_event.is_set(): print(f"[STOP] VK заблокировал аккаунт — прерываю очередь группы {group_id}.") break if not _daily_can_upload(): print(f"[DAILY-LIMIT] Лимит исчерпан — прерываю очередь группы {group_id}.") break # Пауза анти-спама ПЕРЕД постановкой следующего submit (а не # после, как было раньше). Так не висим на пустом sleep после # последнего видео группы и не дублируем интервал, накладывая # его на время самой загрузки. _wait_anti_spam_pause() if _block_event.is_set(): break print(f"\n=== [{group_id}] Планируем видео {idx}/{total_for_group} ===") executor.submit( upload_and_mark, video_path, group_id, idx, total_for_group, state, ) _last_submit_ts = time.time() executor.shutdown(wait=True) # финальная зачистка (fallback) with state_lock: cleanup_completed_files(state) # Если зафиксирована блокировка — выходим со специальным кодом. # Контроллер по rc=42 разворачивает уведомление в Telegram и останавливает # циклы, не делая бесполезных ретраев. if _block_event.is_set(): print("\n[EXIT] VK блокировка → exit code 42 для контроллера.") sys.exit(EXIT_VK_BLOCKED) # Если упёрлись в дневной лимит — отдельный код (не ошибка, а graceful stop). if _effective_daily_limit() > 0 and not _daily_can_upload(): print(f"\n[EXIT] Дневной лимит {_effective_daily_limit()} достигнут → exit code 43.") sys.exit(EXIT_VK_DAILY_LIMIT) print("\nГотово!") if AUTO_DELETE: print("Видео и обложки удаляются только после загрузки во ВСЕ группы (учитывая прошлые запуски).") print("Если выключили свет — при новом запуске продолжит с места остановки (по state).") if __name__ == "__main__": main()
Перезапустить batch_controller после push
Push на выбранные серверы