305 lines
11 KiB
Python
305 lines
11 KiB
Python
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
import re
|
|
import time
|
|
from dataclasses import dataclass
|
|
from typing import Any
|
|
|
|
import aiohttp
|
|
from loguru import logger
|
|
|
|
from .config import settings
|
|
|
|
|
|
class VKAPIError(RuntimeError):
|
|
def __init__(self, code: int | None, message: str) -> None:
|
|
self.code = code
|
|
super().__init__(f"VK API error {code}: {message}")
|
|
|
|
|
|
class VKRateLimiter:
|
|
def __init__(self, rps: int = 3) -> None:
|
|
self.rps = max(1, int(rps))
|
|
self.interval = 1.0 / self.rps
|
|
self._last = 0.0
|
|
self._lock = asyncio.Lock()
|
|
|
|
async def acquire(self) -> None:
|
|
async with self._lock:
|
|
now = time.monotonic()
|
|
wait = self.interval - (now - self._last)
|
|
if wait > 0:
|
|
await asyncio.sleep(wait)
|
|
self._last = time.monotonic()
|
|
|
|
|
|
class VKAPIClient:
|
|
base_url = "https://api.vk.com/method"
|
|
|
|
def __init__(
|
|
self,
|
|
token: str | None = None,
|
|
version: str | None = None,
|
|
rps: int = 3,
|
|
timeout_total_sec: int = 60,
|
|
timeout_connect_sec: int = 10,
|
|
rate_limit_sleep_sec: float = 1.0,
|
|
retry_attempts: int = 3,
|
|
retry_min_delay_sec: float = 2.0,
|
|
retry_max_delay_sec: float = 10.0,
|
|
) -> None:
|
|
self.token = token or settings.vk_access_token
|
|
self.version = version or settings.vk_api_version
|
|
self.limiter = VKRateLimiter(rps)
|
|
self.timeout_total_sec = max(1, int(timeout_total_sec))
|
|
self.timeout_connect_sec = max(1, int(timeout_connect_sec))
|
|
self.rate_limit_sleep_sec = max(0.1, float(rate_limit_sleep_sec))
|
|
self.retry_attempts = max(1, int(retry_attempts))
|
|
self.retry_min_delay_sec = max(0.1, float(retry_min_delay_sec))
|
|
self.retry_max_delay_sec = max(self.retry_min_delay_sec, float(retry_max_delay_sec))
|
|
self.session: aiohttp.ClientSession | None = None
|
|
|
|
async def __aenter__(self) -> "VKAPIClient":
|
|
self.session = aiohttp.ClientSession(
|
|
timeout=aiohttp.ClientTimeout(total=self.timeout_total_sec, connect=self.timeout_connect_sec)
|
|
)
|
|
return self
|
|
|
|
async def __aexit__(self, *args) -> None:
|
|
if self.session:
|
|
await self.session.close()
|
|
|
|
async def call(self, method: str, **params: Any) -> dict | list:
|
|
if not self.session:
|
|
raise RuntimeError("VKAPIClient is not initialized")
|
|
if not self.token:
|
|
raise RuntimeError("VK_ACCESS_TOKEN is empty")
|
|
|
|
payload = dict(params)
|
|
payload["access_token"] = self.token
|
|
payload["v"] = self.version
|
|
|
|
for attempt in range(1, self.retry_attempts + 1):
|
|
try:
|
|
await self.limiter.acquire()
|
|
async with self.session.post(f"{self.base_url}/{method}", data=payload) as resp:
|
|
resp.raise_for_status()
|
|
data = await resp.json(content_type=None)
|
|
if "error" not in data:
|
|
return data.get("response", {})
|
|
|
|
err = data["error"]
|
|
code = err.get("error_code")
|
|
message = err.get("error_msg", "unknown")
|
|
if code == 6 and attempt < self.retry_attempts:
|
|
logger.warning("VK rate limit hit, sleeping {}s", self.rate_limit_sleep_sec)
|
|
await asyncio.sleep(self.rate_limit_sleep_sec)
|
|
continue
|
|
raise VKAPIError(code, message)
|
|
except VKAPIError:
|
|
raise
|
|
except Exception:
|
|
if attempt >= self.retry_attempts:
|
|
raise
|
|
delay = min(self.retry_min_delay_sec * (2 ** (attempt - 1)), self.retry_max_delay_sec)
|
|
await asyncio.sleep(delay)
|
|
|
|
raise VKAPIError(None, "retry exhausted")
|
|
|
|
async def resolve_group(self, input_value: str) -> tuple[str, int, str]:
|
|
external_id = normalize_vk_source(input_value)
|
|
if external_id.lstrip("-").isdigit():
|
|
group_id = abs(int(external_id))
|
|
response = await self.call("groups.getById", group_id=str(group_id))
|
|
else:
|
|
response = await self.call("groups.getById", group_id=external_id)
|
|
items = response if isinstance(response, list) else response.get("groups", [])
|
|
if not items:
|
|
raise VKAPIError(None, f"cannot resolve VK group: {input_value}")
|
|
group = items[0]
|
|
group_id = int(group["id"])
|
|
screen_name = str(group.get("screen_name") or external_id)
|
|
name = str(group.get("name") or screen_name)
|
|
return screen_name, -group_id, name
|
|
|
|
async def get_wall_posts(self, owner_id: int, count: int, offset: int = 0) -> dict:
|
|
response = await self.call("wall.get", owner_id=owner_id, count=count, offset=offset, filter="owner")
|
|
if not isinstance(response, dict):
|
|
raise VKAPIError(None, "wall.get returned non-object response")
|
|
return response
|
|
|
|
async def get_wall_upload_server(self, group_id: int) -> str:
|
|
response = await self.call("photos.getWallUploadServer", group_id=abs(int(group_id)))
|
|
if not isinstance(response, dict) or not response.get("upload_url"):
|
|
raise VKAPIError(None, "photos.getWallUploadServer returned no upload_url")
|
|
return str(response["upload_url"])
|
|
|
|
async def upload_wall_photo_bytes(self, group_id: int, data: bytes, filename: str = "photo.jpg") -> dict:
|
|
if not self.session:
|
|
raise RuntimeError("VKAPIClient is not initialized")
|
|
upload_url = await self.get_wall_upload_server(group_id)
|
|
form = aiohttp.FormData()
|
|
form.add_field("photo", data, filename=filename, content_type="image/jpeg")
|
|
async with self.session.post(upload_url, data=form) as resp:
|
|
uploaded = await resp.json(content_type=None)
|
|
saved = await self.call(
|
|
"photos.saveWallPhoto",
|
|
group_id=abs(int(group_id)),
|
|
photo=uploaded.get("photo"),
|
|
server=uploaded.get("server"),
|
|
hash=uploaded.get("hash"),
|
|
)
|
|
if not isinstance(saved, list) or not saved:
|
|
raise VKAPIError(None, f"photos.saveWallPhoto returned invalid response: {json.dumps(saved)[:300]}")
|
|
return saved[0]
|
|
|
|
async def upload_wall_photo_url(self, group_id: int, url: str) -> dict:
|
|
if not self.session:
|
|
raise RuntimeError("VKAPIClient is not initialized")
|
|
async with self.session.get(url) as resp:
|
|
resp.raise_for_status()
|
|
content_type = resp.headers.get("content-type") or "image/jpeg"
|
|
data = await resp.read()
|
|
suffix = ".jpg"
|
|
if "png" in content_type:
|
|
suffix = ".png"
|
|
elif "webp" in content_type:
|
|
suffix = ".webp"
|
|
return await self.upload_wall_photo_bytes(group_id, data, filename=f"photo{suffix}")
|
|
|
|
async def create_wall_post(
|
|
self,
|
|
owner_id: int,
|
|
message: str,
|
|
attachments: list[str],
|
|
from_group: bool = True,
|
|
) -> int:
|
|
response = await self.call(
|
|
"wall.post",
|
|
owner_id=int(owner_id),
|
|
from_group=1 if from_group else 0,
|
|
message=message or "",
|
|
attachments=",".join(attachments) if attachments else "",
|
|
)
|
|
if not isinstance(response, dict) or "post_id" not in response:
|
|
raise VKAPIError(None, f"wall.post returned invalid response: {response}")
|
|
return int(response["post_id"])
|
|
|
|
|
|
def normalize_vk_source(value: str) -> str:
|
|
raw = str(value or "").strip()
|
|
if not raw:
|
|
return ""
|
|
if "vk.com" in raw:
|
|
raw = raw.split("vk.com", 1)[1]
|
|
raw = raw.lstrip("/")
|
|
raw = raw.split("?", 1)[0].split("#", 1)[0].strip()
|
|
raw = re.sub(r"^(club|public)", "", raw, flags=re.IGNORECASE)
|
|
return raw.lower()
|
|
|
|
|
|
def is_repost(post: dict) -> bool:
|
|
return bool(post.get("copy_history"))
|
|
|
|
|
|
def is_deleted_or_invalid(post: dict) -> bool:
|
|
if post.get("is_deleted") or post.get("deleted"):
|
|
return True
|
|
if not post.get("id") or not post.get("date"):
|
|
return True
|
|
return False
|
|
|
|
|
|
def is_fatal_source_error(error: Exception | str) -> bool:
|
|
message = str(error or "").lower()
|
|
fatal_patterns = (
|
|
"vk api error 15:",
|
|
"vk api error 18:",
|
|
"vk api error 19:",
|
|
"vk api error 100:",
|
|
"vk api error 113:",
|
|
"vk api error 1051:",
|
|
"vk api error 200:",
|
|
"vk api error 201:",
|
|
"vk api error 203:",
|
|
"[15]",
|
|
"[18]",
|
|
"[19]",
|
|
"[100]",
|
|
"[113]",
|
|
"[1051]",
|
|
"[200]",
|
|
"[201]",
|
|
"[203]",
|
|
"access denied",
|
|
"private profile",
|
|
"cannot resolve vk group",
|
|
"cannot resolve screen_name",
|
|
"method is unavailable with current profile type",
|
|
"wall is disabled",
|
|
)
|
|
return any(pattern in message for pattern in fatal_patterns)
|
|
|
|
|
|
def post_vk_url(owner_id: int, post_id: int | str) -> str:
|
|
return f"https://vk.com/wall{int(owner_id)}_{post_id}"
|
|
|
|
|
|
@dataclass
|
|
class ExtractedMedia:
|
|
media_type: str
|
|
original_url: str | None
|
|
attachment_id: str
|
|
width: int | None = None
|
|
height: int | None = None
|
|
duration_sec: int | None = None
|
|
sort_order: int = 0
|
|
|
|
|
|
def extract_media(post: dict) -> list[ExtractedMedia]:
|
|
out: list[ExtractedMedia] = []
|
|
for idx, att in enumerate(post.get("attachments", []) or []):
|
|
att_type = att.get("type")
|
|
if att_type == "photo":
|
|
photo = att.get("photo") or {}
|
|
sizes = sorted(
|
|
photo.get("sizes", []) or [],
|
|
key=lambda s: int(s.get("width") or 0) * int(s.get("height") or 0),
|
|
reverse=True,
|
|
)
|
|
best = sizes[0] if sizes else {}
|
|
owner_id = photo.get("owner_id")
|
|
media_id = photo.get("id")
|
|
if owner_id is None or media_id is None:
|
|
continue
|
|
out.append(
|
|
ExtractedMedia(
|
|
media_type="photo",
|
|
original_url=best.get("url"),
|
|
attachment_id=f"photo{owner_id}_{media_id}",
|
|
width=best.get("width"),
|
|
height=best.get("height"),
|
|
sort_order=idx,
|
|
)
|
|
)
|
|
elif att_type == "video":
|
|
video = att.get("video") or {}
|
|
owner_id = video.get("owner_id")
|
|
media_id = video.get("id")
|
|
if owner_id is None or media_id is None:
|
|
continue
|
|
out.append(
|
|
ExtractedMedia(
|
|
media_type="video",
|
|
original_url=f"https://vk.com/video{owner_id}_{media_id}",
|
|
attachment_id=f"video{owner_id}_{media_id}",
|
|
width=video.get("width"),
|
|
height=video.get("height"),
|
|
duration_sec=video.get("duration"),
|
|
sort_order=idx,
|
|
)
|
|
)
|
|
return out
|