oss.py 4.24 KB
import requests
import oss2
import math
import time
from urllib.parse import urlparse
from requests.adapters import HTTPAdapter
from .config import (
    HTTP_POOL_MAXSIZE, HTTP_TRANSFER_RETRIES, HTTP_TRANSFER_RETRY_BACKOFF_SECONDS, OSS_CONFIG,
)
from .utils import compute_audio_md5

_HTTP_SESSION = requests.Session()
_HTTP_ADAPTER = HTTPAdapter(pool_connections=HTTP_POOL_MAXSIZE, pool_maxsize=HTTP_POOL_MAXSIZE)
_HTTP_SESSION.mount('http://', _HTTP_ADAPTER)
_HTTP_SESSION.mount('https://', _HTTP_ADAPTER)


def _http_get(url: str, timeout: int = 30):
    return _HTTP_SESSION.get(url, timeout=timeout)


def _clean_url(url) -> str:
    if url is None:
        return ''
    if isinstance(url, float) and math.isnan(url):
        return ''
    text = str(url).strip()
    if not text or text.lower() in {'nan', 'none', 'null'}:
        return ''
    parsed = urlparse(text)
    if parsed.scheme not in {'http', 'https'} or not parsed.netloc:
        return ''
    return text


def _host(url: str | None) -> str:
    return urlparse(url or '').netloc.lower()


def _is_target_oss_url(url: str | None, base_url: str) -> bool:
    return bool(url) and _host(url) == _host(base_url)


def _download_url(url: str, base_url: str) -> str:
    """Optionally rewrite target-bucket downloads to an internal/CNAME base URL."""
    download_base_url = OSS_CONFIG.get('download_base_url')
    rewrite_from = OSS_CONFIG.get('download_rewrite_from_base_url') or base_url
    if not download_base_url or not url.startswith(rewrite_from.rstrip('/') + '/'):
        return url
    return f"{download_base_url.rstrip('/')}{url[len(rewrite_from.rstrip('/')):]}"


def _download_content(url: str, base_url: str) -> bytes:
    """下载资源;针对超时、连接中断和服务端错误进行有限重试。"""
    download_url = _download_url(url, base_url)
    attempts = max(1, HTTP_TRANSFER_RETRIES)
    for attempt in range(attempts):
        try:
            resp = _http_get(download_url, timeout=30)
            resp.raise_for_status()
            return resp.content
        except requests.RequestException as exc:
            status_code = getattr(getattr(exc, 'response', None), 'status_code', None)
            # 4xx(429 除外)是确定性失败,不浪费时间重试。
            if status_code is not None and 400 <= status_code < 500 and status_code != 429:
                raise
            if attempt == attempts - 1:
                raise
            time.sleep(HTTP_TRANSFER_RETRY_BACKOFF_SECONDS * (2 ** attempt))


def transfer_url(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: str) -> str:
    """
    将 url 指向的文件转移到 archive-dev OSS 的 oss_key 路径。
    若 url 为空或已在目标 bucket,直接返回原 url(不上传)。
    返回新的公开访问 URL。
    """
    url = _clean_url(url)
    if not url:
        return ''
    if _is_target_oss_url(url, base_url):
        return url
    content = _download_content(url, base_url)
    bucket.put_object(oss_key, content)
    return f"{base_url.rstrip('/')}/{oss_key}"


def download_text_url(url: str | None, base_url: str) -> str:
    """下载歌词等文本 URL;即使源文件已在目标 OSS,也实际读取其内容。"""
    url = _clean_url(url)
    if not url:
        return ''
    content = _download_content(url, base_url)
    return content.decode('utf-8-sig')


def transfer_url_with_md5(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: str) -> tuple[str, str]:
    """
    将音频 URL 转存到 OSS,并基于下载到的音频字节计算 MD5。
    已在目标 OSS 的 URL 无需重新下载,无法可靠计算 MD5,返回空 MD5。
    """
    url = _clean_url(url)
    if not url:
        return '', ''
    if _is_target_oss_url(url, base_url):
        return url, ''
    content = _download_content(url, base_url)
    audio_md5 = compute_audio_md5(content)
    bucket.put_object(oss_key, content)
    return f"{base_url.rstrip('/')}/{oss_key}", audio_md5


def build_oss_key(platform: str, category: str, filename: str) -> str:
    """
    构造 archive-dev 内的存储路径。
    platform: 'qq' | 'kugou' | 'netease'
    category: 'audio' | 'cover' | 'singer' | 'album'
    filename: 带扩展名的文件名
    """
    return f"crawler/{category}/{platform}/{filename}"