Files
2026-07-17 09:24:47 +08:00

1042 lines
38 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""抖音 IM 私信图片上传。
复刻抖音 PC 网页版私信发图的真实链路(与"创作者中心 ImageX 图文图床"是两套,
IM 发送只认这条链路产出的 tos-cn-o-* 资源):
1. GET www.douyin.com/aweme/v1/web/im/upload/config/v2 a_bogus 签名)
→ 返回内含 STS2 凭证(AccessKeyID + SignedSecretAccessKey)与 SpaceName(=zhenzhen)
2. GET vod.bytedanceapi.com/?Action=ApplyUploadInner&SpaceName=zhenzhen&FileType=image
AWS4-HMAC-SHA256service=vod,带 x-amz-security-token=STS2…)
→ 返回 UploadHost / StoreUri(tos-cn-o-*) / Auth(SpaceKey JWT) / SessionKey
3. POST https://{UploadHost}/upload/v1/{StoreUri}
Authorization: SpaceKey/zhenzhen/…JWTContent-CRC32X-Storage-U=my_uid
4. POST vod.bytedanceapi.com/?Action=CommitUploadInner&SpaceName=zhenzhen AWS4 签名)
→ 确认上传,最终 uri 即 StoreUritos-cn-o-*
之后用该 tos-cn-o-* uri 构造 type=27 消息体即可被 IM 后端校验通过。
"""
from __future__ import annotations
import base64
import datetime
import hashlib
import hmac
import json
import logging
import random
import re
import string
import zlib
from typing import Any
from urllib.parse import urlencode
logger = logging.getLogger("douyin_im.image_upload")
_LOCAL_URL_RE = re.compile(
r"^(/api/media/(?:messages|link-cards)/|"
r"https?://(?:localhost|127\.0\.0\.1)(?::\d+)?/api/media/(?:messages|link-cards)/)",
re.I,
)
# 抖音 web 私信图片上传专用通道(VOD/ImageX inner 接口 + IM 上传配置)
VOD_HOST = "https://vod.bytedanceapi.com/"
VOD_REGION = "cn-north-1"
VOD_SERVICE = "vod"
IM_UPLOAD_CONFIG_URL = "https://www.douyin.com/aweme/v1/web/im/upload/config/v2"
DEFAULT_SPACE_NAME = "zhenzhen"
def _public_base_url() -> str:
"""读取服务器对外域名(与 link_cards.public_base_url 一致),不依赖 Request 对象。"""
import os
env = os.getenv("KEFU_PUBLIC_BASE_URL", "").strip().rstrip("/")
return env
def is_local_media_url(url: str) -> bool:
raw = (url or "").strip()
if not raw:
return False
if raw.startswith("/api/media/messages/"):
return True
if raw.startswith("/api/media/link-cards/"):
return True
if _LOCAL_URL_RE.match(raw):
return True
# 部署在真实域名时(配置了 KEFU_PUBLIC_BASE_URL),cover_url 可能是 http(s)://真实域名/api/media/...
public_base = _public_base_url()
if public_base and raw.startswith(public_base):
path_part = raw[len(public_base):]
if path_part.startswith("/api/media/messages/") or path_part.startswith("/api/media/link-cards/"):
return True
return False
def is_douyin_cdn_url(url: str) -> bool:
raw = (url or "").strip().lower()
if not raw.startswith("http"):
return False
return any(
host in raw
for host in (
"douyinpic.com",
"byteimg.com",
"ibyteimg.com",
"douyinstatic.com",
"snssdk.com",
"vodupload.com",
)
)
def parse_local_message_media_path(url: str) -> tuple[int, str] | None:
"""从 /api/media/messages/{account_id}/{filename} 解析 account_id 与文件名。"""
raw = (url or "").strip()
m = re.search(r"/api/media/messages/(\d+)/([^/?#]+)", raw)
if not m:
return None
return int(m.group(1)), m.group(2)
def normalize_cover_cache_key(cover_ref: str) -> str:
"""本地封面路径的稳定缓存键(避免同一文件因域名前缀不同重复上传)。"""
raw = (cover_ref or "").strip().split("?")[0]
if not raw:
return ""
parsed = parse_link_card_media_path(raw)
if parsed:
owner_id, fname = parsed
return f"/api/media/link-cards/{owner_id}/{fname}"
msg_parsed = parse_local_message_media_path(raw)
if msg_parsed:
account_id, fname = msg_parsed
return f"/api/media/messages/{account_id}/{fname}"
if raw.startswith("/api/media/"):
return raw
return raw
def _cover_cache_file(upload_dir: str) -> str:
import os
return os.path.join(os.path.dirname(upload_dir), "im_cover_uri_cache.json")
def load_cached_cover_upload(cover_ref: str, upload_dir: str) -> dict[str, Any] | None:
"""读取已上传封面缓存,避免每次发卡片都重新上传拿到未过审的新 uri。"""
import os
key = normalize_cover_cache_key(cover_ref)
if not key:
return None
path = _cover_cache_file(upload_dir)
if not os.path.isfile(path):
return None
try:
with open(path, encoding="utf-8") as f:
cache = json.load(f)
if not isinstance(cache, dict):
return None
entry = cache.get(key)
if isinstance(entry, dict) and entry.get("uri"):
return entry
except Exception as exc:
logger.debug("load_cached_cover_upload failed: %s", exc)
return None
def save_cached_cover_upload(cover_ref: str, upload_dir: str, uploaded: dict[str, Any]) -> None:
"""持久化封面 uri,供后续链接卡片/图片发送复用。"""
import os
key = normalize_cover_cache_key(cover_ref)
uri = str(uploaded.get("uri") or "").strip()
if not key or not uri:
return
path = _cover_cache_file(upload_dir)
cache: dict[str, Any] = {}
if os.path.isfile(path):
try:
with open(path, encoding="utf-8") as f:
loaded = json.load(f)
if isinstance(loaded, dict):
cache = loaded
except Exception:
cache = {}
cache[key] = {
"uri": uri,
"url": str(uploaded.get("url") or "").strip(),
"md5": str(uploaded.get("md5") or "").strip(),
"url_list": uploaded.get("url_list") if isinstance(uploaded.get("url_list"), list) else [],
}
try:
os.makedirs(os.path.dirname(path), exist_ok=True)
with open(path, "w", encoding="utf-8") as f:
json.dump(cache, f, ensure_ascii=False, indent=2)
except Exception as exc:
logger.warning("save_cached_cover_upload failed: %s", exc)
def _apply_cover_upload_to_spec(spec: dict[str, Any], uploaded: dict[str, Any]) -> None:
if uploaded.get("url"):
spec["cover_url"] = uploaded["url"]
if uploaded.get("uri"):
spec["cover_uri"] = uploaded["uri"]
if uploaded.get("md5"):
spec["cover_md5"] = uploaded["md5"]
uri = str(uploaded.get("uri") or spec.get("cover_uri") or "").strip().lstrip("/")
if uri:
from .message_content import uri_to_cdn_urls
urls = uploaded.get("url_list")
if isinstance(urls, list) and urls:
spec["url_list"] = [str(u) for u in urls if str(u).strip()]
else:
spec["url_list"] = uri_to_cdn_urls(uri)
def _fill_card_page_url(spec: dict[str, Any]) -> None:
page_url = str(spec.get("url") or spec.get("page_url") or "").strip()
if page_url.startswith("/"):
base = _public_base_url().rstrip("/")
if base:
spec["url"] = f"{base}{page_url}"
elif not page_url:
target = str(spec.get("target_url") or spec.get("link_url") or "").strip()
if target:
spec["url"] = target if target.startswith("http") else f"https://{target.lstrip('/')}"
def parse_link_card_media_path(url: str) -> tuple[int, str] | None:
"""从 /api/media/link-cards/{owner_id}/{filename} 解析 owner_id 与文件名。"""
raw = (url or "").strip()
m = re.search(r"/api/media/link-cards/(\d+)/([^/?#]+)", raw)
if not m:
return None
return int(m.group(1)), m.group(2)
def _requests_proxies() -> dict | None:
try:
from rpa_engine.runtime_config import requests_proxies
return requests_proxies()
except Exception:
return None
def _random_s() -> str:
chars = string.digits + string.ascii_lowercase
return "".join(random.choice(chars) for _ in range(11))
def _signing_key(secret: str, date_stamp: str, region: str, service: str) -> bytes:
k_date = hmac.new(("AWS4" + secret).encode(), date_stamp.encode(), hashlib.sha256).digest()
k_region = hmac.new(k_date, region.encode(), hashlib.sha256).digest()
k_service = hmac.new(k_region, service.encode(), hashlib.sha256).digest()
return hmac.new(k_service, b"aws4_request", hashlib.sha256).digest()
def _aws4_authorization(
*,
method: str,
canonical_querystring: str,
amz_date: str,
date_stamp: str,
session_token: str,
access_key_id: str,
secret_access_key: str,
payload_hash: str | None = None,
signed_headers: list[str] | None = None,
service: str = VOD_SERVICE,
region: str = VOD_REGION,
) -> str:
if signed_headers is None:
signed_headers = ["x-amz-date", "x-amz-security-token"]
if payload_hash is None:
payload_hash = hashlib.sha256(b"").hexdigest()
header_lines = []
for name in signed_headers:
if name == "x-amz-date":
header_lines.append(f"x-amz-date:{amz_date}\n")
elif name == "x-amz-security-token":
header_lines.append(f"x-amz-security-token:{session_token}\n")
elif name == "x-amz-content-sha256":
header_lines.append(f"x-amz-content-sha256:{payload_hash}\n")
canonical_headers = "".join(header_lines)
signed = ";".join(signed_headers)
canonical_request = (
f"{method}\n/\n{canonical_querystring}\n{canonical_headers}\n{signed}\n{payload_hash}"
)
credential_scope = f"{date_stamp}/{region}/{service}/aws4_request"
string_to_sign = (
"AWS4-HMAC-SHA256\n"
f"{amz_date}\n"
f"{credential_scope}\n"
f"{hashlib.sha256(canonical_request.encode()).hexdigest()}"
)
signature = hmac.new(
_signing_key(secret_access_key, date_stamp, region, service),
string_to_sign.encode(),
hashlib.sha256,
).hexdigest()
return (
f"AWS4-HMAC-SHA256 Credential={access_key_id}/{credential_scope}, "
f"SignedHeaders={signed}, Signature={signature}"
)
def _safe_json(resp) -> dict[str, Any]:
text = (getattr(resp, "text", None) or "").strip()
if not text:
return {"error": f"空响应 (HTTP {getattr(resp, 'status_code', '?')})"}
try:
data = resp.json()
return data if isinstance(data, dict) else {"error": "响应不是 JSON 对象"}
except Exception:
snippet = text[:200].replace("\n", " ")
return {"error": f"非 JSON 响应 (HTTP {getattr(resp, 'status_code', '?')}): {snippet}"}
# ---------------------------------------------------------------------------
# 第 1 步:拉取 IM 上传配置,提取 STS 凭证 + SpaceName
# ---------------------------------------------------------------------------
def _sanitize_for_log(obj: Any, depth: int = 0) -> Any:
"""结构化脱敏:保留字段名与层级,长字符串只留前 12 位 + 长度,便于排查凭证字段。"""
if depth > 6:
return "…"
if isinstance(obj, dict):
return {str(k): _sanitize_for_log(v, depth + 1) for k, v in obj.items()}
if isinstance(obj, list):
return [_sanitize_for_log(v, depth + 1) for v in obj[:3]]
if isinstance(obj, str):
return f"{obj[:12]}…(len={len(obj)})" if len(obj) > 24 else obj
return obj
def _find_auth_object(obj: Any, depth: int = 0) -> dict | None:
"""找到「直接含有 STS2 会话凭证字符串」的那个 dict(即临时凭证对象)。
该对象里通常同时含有 AccessKeyId / SecretAccessKey / SessionToken(=STS2…)。
关键:签名要用同级的 SecretAccessKey 字段,而**不是** STS2 令牌内部解出的
SignedSecretAccessKey(那是服务端校验用的,拿来当签名密钥会 SignatureDoesNotMatch)。
"""
if depth > 8 or not isinstance(obj, (dict, list)):
return None
if isinstance(obj, dict):
for value in obj.values():
if isinstance(value, str) and value.startswith("STS2"):
return obj
for value in obj.values():
found = _find_auth_object(value, depth + 1)
if found is not None:
return found
else:
for value in obj:
found = _find_auth_object(value, depth + 1)
if found is not None:
return found
return None
def _find_sts_token(obj: Any, depth: int = 0) -> str:
"""递归在响应里找到形如 'STS2...' 的会话凭证字符串。"""
if depth > 8:
return ""
if isinstance(obj, str):
return obj if obj.startswith("STS2") else ""
if isinstance(obj, dict):
for value in obj.values():
found = _find_sts_token(value, depth + 1)
if found:
return found
elif isinstance(obj, list):
for value in obj:
found = _find_sts_token(value, depth + 1)
if found:
return found
return ""
def _find_space_name(obj: Any, depth: int = 0) -> str:
if depth > 8:
return ""
if isinstance(obj, dict):
for key, value in obj.items():
if str(key).lower() in ("space_name", "spacename") and isinstance(value, str) and value:
return value
found = _find_space_name(value, depth + 1)
if found:
return found
elif isinstance(obj, list):
for value in obj:
found = _find_space_name(value, depth + 1)
if found:
return found
return ""
def _decode_sts(sts_token: str) -> tuple[str, str]:
"""STS2<base64(JSON)>,解出 AccessKeyID 与 SignedSecretAccessKey。"""
try:
b64 = sts_token[4:] if sts_token.startswith("STS2") else sts_token
b64 += "=" * (-len(b64) % 4)
data = json.loads(base64.b64decode(b64).decode("utf-8", "ignore"))
ak = data.get("AccessKeyID") or data.get("AccessKeyId") or ""
sk = data.get("SignedSecretAccessKey") or data.get("SecretAccessKey") or ""
return ak, sk
except Exception:
return "", ""
def _fetch_im_upload_sts(session) -> tuple[str, str, str, str]:
"""返回 (access_key_id, secret_access_key, sts_token, space_name)。"""
import requests
from .auth import DouyinAuth
from .dy_util import (
DEFAULT_USER_AGENT,
generate_a_bogus,
generate_msToken,
generate_webid,
splice_url,
)
auth = DouyinAuth()
auth.perepare_auth(session.cookie_header(), session.web_protect_str, session.keys_str)
ua = session.user_agent or DEFAULT_USER_AGENT
params = {
"device_platform": "webapp",
"aid": "6383",
"channel": "channel_pc_web",
"update_version_code": "170400",
"pc_client_type": "1",
"pc_libra_divert": "Windows",
"support_h265": "1",
"support_dash": "1",
"version_code": "170400",
"version_name": "17.4.0",
"cookie_enabled": "true",
"screen_width": "1536",
"screen_height": "960",
"browser_language": "zh-CN",
"browser_platform": "Win32",
"browser_name": "Chrome",
"browser_version": "120.0.0.0",
"browser_online": "true",
"engine_name": "Blink",
"engine_version": "120.0.0.0",
"os_name": "Windows",
"os_version": "10",
"cpu_core_num": "8",
"device_memory": "8",
"platform": "PC",
"downlink": "10",
"effective_type": "4g",
"round_trip_time": "50",
"webid": generate_webid(auth, "https://www.douyin.com/"),
"verifyFp": auth.cookie.get("s_v_web_id", "") if auth.cookie else "",
"fp": auth.cookie.get("s_v_web_id", "") if auth.cookie else "",
"msToken": auth.msToken or generate_msToken(),
}
query = splice_url(params)
params["a_bogus"] = generate_a_bogus(query, user_agent=ua)
headers = {
"User-Agent": ua,
"Referer": "https://www.douyin.com/",
"Accept": "application/json, text/plain, */*",
}
resp = requests.get(
IM_UPLOAD_CONFIG_URL,
params=params,
headers=headers,
cookies=auth.cookie,
timeout=20,
verify=False,
proxies=_requests_proxies(),
)
data = _safe_json(resp)
if data.get("error"):
raise RuntimeError(f"获取 IM 上传配置失败:{data['error']}")
if data.get("status_code") not in (None, 0):
raise RuntimeError(
f"获取 IM 上传配置失败:status_code={data.get('status_code')} "
f"{data.get('status_msg') or ''}"
)
# 诊断:把 config/v2 响应结构(字段名保留、长字符串脱敏)打到日志,便于核对凭证字段。
try:
logger.info(
"im/upload/config/v2 结构: %s",
json.dumps(_sanitize_for_log(data), ensure_ascii=False),
)
except Exception:
pass
# 凭证块选择(决定 8003 是否发生的关键):
# - inner_image_config / public_image_config / public_file_config 共用同一套
# 可用于 VOD 上传的 STS 令牌,仅 space 不同;
# - public_image_config_v2 的 token 不同,不兼容 VOD 内部上传(会报
# "session token sequence is broken"),必须避开。
# 空间含义:maya_review = 审核空间,图片上传后处于待审核态,作为消息发送会被拒(8003);
# zhenzhen = 公开可发送空间,IM 图片消息应使用它。
# 因此优先 public_image_configzhenzhen),其凭证同样能完成 VOD 上传。
def _block_has_sts(block: Any) -> bool:
return isinstance(block, dict) and any(
isinstance(v, str) and v.startswith("STS2") for v in block.values()
)
auth_obj = None
for _key in ("public_image_config", "inner_image_config"):
if isinstance(data, dict) and _block_has_sts(data.get(_key)):
auth_obj = data[_key]
break
if auth_obj is None:
auth_obj = _find_auth_object(data)
if not auth_obj:
raise RuntimeError(
"IM 上传配置响应里未找到 STS 凭证(STS2 token);可能 cookie/签名失效,请用浏览器模式重新登录"
)
sts = next(
(v for v in auth_obj.values() if isinstance(v, str) and v.startswith("STS2")),
"",
)
# 优先用凭证对象同级的 AccessKeyId / SecretAccessKey(用于 SigV4 签名的真实密钥)。
ak = (
auth_obj.get("AccessKeyID")
or auth_obj.get("AccessKeyId")
or auth_obj.get("access_key_id")
or ""
)
sk = (
auth_obj.get("SecretAccessKey")
or auth_obj.get("SecretAccesskey")
or auth_obj.get("secret_access_key")
or ""
)
# 兜底:若响应没给独立的 ak/sk,再尝试从 STS2 令牌解码(SignedSecretAccessKey 一般不可用,仅最后兜底)。
if not ak or not sk:
dec_ak, dec_sk = _decode_sts(sts)
ak = ak or dec_ak
sk = sk or dec_sk
if not ak or not sk:
raise RuntimeError("解析 STS 凭证失败(缺少 AccessKeyId / SecretAccessKey")
# space 必须取自与凭证同一块,避免凭证用 inner_image_config 而 space 误取到别处。
space = (
auth_obj.get("space_name")
or auth_obj.get("SpaceName")
or auth_obj.get("spaceName")
or _find_space_name(data)
or DEFAULT_SPACE_NAME
)
return ak, sk, sts, space
# ---------------------------------------------------------------------------
# 第 2 步:ApplyUploadInnerVOD),申请上传地址
# ---------------------------------------------------------------------------
def _extract_apply_inner(data: dict[str, Any]) -> tuple[str, str, str, str]:
"""返回 (upload_host, store_uri, jwt_auth, session_key)。"""
result = data.get("Result") or {}
addr = result.get("InnerUploadAddress") or result.get("UploadAddress") or {}
nodes = addr.get("UploadNodes") or []
if nodes:
node = nodes[0]
stores = node.get("StoreInfos") or []
store = stores[0] if stores else {}
host = node.get("UploadHost") or ""
if not host:
hosts = node.get("UploadHosts") or addr.get("UploadHosts") or []
host = hosts[0] if hosts else ""
return (
str(host or ""),
str(store.get("StoreUri") or ""),
str(store.get("Auth") or ""),
str(node.get("SessionKey") or addr.get("SessionKey") or ""),
)
hosts = addr.get("UploadHosts") or []
host = hosts[0] if hosts else ""
stores = addr.get("StoreInfos") or []
store = stores[0] if stores else {}
return (
str(host or ""),
str(store.get("StoreUri") or ""),
str(store.get("Auth") or ""),
str(addr.get("SessionKey") or result.get("SessionKey") or ""),
)
def _vod_apply_upload_inner(
ak: str, sk: str, token: str, space: str, file_size: int
) -> tuple[str, str, str, str]:
import requests
from .dy_util import DEFAULT_USER_AGENT
now = datetime.datetime.utcnow()
amz_date = now.strftime("%Y%m%dT%H%M%SZ")
date_stamp = now.strftime("%Y%m%d")
params = {
"Action": "ApplyUploadInner",
"Version": "2020-11-19",
"SpaceName": space,
"FileType": "image",
"IsInner": "1",
"NeedFallback": "true",
"FileSize": str(file_size),
"s": _random_s(),
}
qs = urlencode(sorted(params.items()))
authorization = _aws4_authorization(
method="GET",
canonical_querystring=qs,
amz_date=amz_date,
date_stamp=date_stamp,
session_token=token,
access_key_id=ak,
secret_access_key=sk,
service=VOD_SERVICE,
)
resp = requests.get(
f"{VOD_HOST}?{qs}",
headers={
"accept": "*/*",
"authorization": authorization,
"user-agent": DEFAULT_USER_AGENT,
"x-amz-date": amz_date,
"x-amz-security-token": token,
"Referer": "https://www.douyin.com/",
},
timeout=30,
verify=False,
proxies=_requests_proxies(),
)
data = _safe_json(resp)
if data.get("error"):
raise RuntimeError(f"申请上传地址失败:{data['error']}")
meta = data.get("ResponseMetadata") or {}
err = meta.get("Error")
if err:
raise RuntimeError(f"申请上传地址失败:{err.get('Message') or err.get('Code')}")
return _extract_apply_inner(data)
# ---------------------------------------------------------------------------
# 第 3 步:上传二进制(SpaceKey JWT 鉴权)
# ---------------------------------------------------------------------------
def _vod_upload_binary(
host: str, store_uri: str, jwt_auth: str, user_id: str, raw: bytes, session=None
) -> None:
import requests
from .dy_util import DEFAULT_USER_AGENT
crc32 = format(zlib.crc32(raw) & 0xFFFFFFFF, "08x")
url = f"https://{host}/upload/v1/{store_uri}"
ua = (getattr(session, "user_agent", None) or DEFAULT_USER_AGENT)
headers = {
"Authorization": jwt_auth,
"Content-CRC32": crc32,
"Content-Type": "application/octet-stream",
"Content-Disposition": 'attachment; filename="undefined"',
"User-Agent": ua,
}
if user_id:
headers["X-Storage-U"] = str(user_id)
resp = requests.post(
url,
headers=headers,
data=raw,
timeout=60,
verify=False,
proxies=_requests_proxies(),
)
data = _safe_json(resp)
if data.get("error"):
raise RuntimeError(f"上传图片数据失败:{data['error']}")
if data.get("code") not in (2000, 0, None):
raise RuntimeError(f"上传图片数据失败:{data.get('message') or data}")
# ---------------------------------------------------------------------------
# 第 4 步:CommitUploadInnerVOD),确认上传
# ---------------------------------------------------------------------------
def _vod_commit_upload_inner(
ak: str, sk: str, token: str, space: str, session_key: str
) -> dict[str, Any]:
import requests
from .dy_util import DEFAULT_USER_AGENT
now = datetime.datetime.utcnow()
amz_date = now.strftime("%Y%m%dT%H%M%SZ")
date_stamp = now.strftime("%Y%m%d")
params = {
"Action": "CommitUploadInner",
"Version": "2020-11-19",
"SpaceName": space,
}
qs = urlencode(sorted(params.items()))
body = json.dumps({"SessionKey": session_key, "Functions": []}, separators=(",", ":"))
payload_hash = hashlib.sha256(body.encode()).hexdigest()
signed_headers = ["x-amz-content-sha256", "x-amz-date", "x-amz-security-token"]
authorization = _aws4_authorization(
method="POST",
canonical_querystring=qs,
amz_date=amz_date,
date_stamp=date_stamp,
session_token=token,
access_key_id=ak,
secret_access_key=sk,
payload_hash=payload_hash,
signed_headers=signed_headers,
service=VOD_SERVICE,
)
resp = requests.post(
f"{VOD_HOST}?{qs}",
data=body,
headers={
"accept": "*/*",
"authorization": authorization,
"content-type": "application/json",
"user-agent": DEFAULT_USER_AGENT,
"x-amz-content-sha256": payload_hash,
"x-amz-date": amz_date,
"x-amz-security-token": token,
"Referer": "https://www.douyin.com/",
},
timeout=30,
verify=False,
proxies=_requests_proxies(),
)
data = _safe_json(resp)
if data.get("error"):
raise RuntimeError(f"确认上传失败:{data['error']}")
meta = data.get("ResponseMetadata") or {}
err = meta.get("Error")
if err:
raise RuntimeError(f"确认上传失败:{err.get('Message') or err.get('Code')}")
return data.get("Result") or {}
def upload_im_image(
session,
raw: bytes,
*,
filename: str = "image.jpg",
content_type: str = "image/jpeg",
) -> dict[str, Any]:
"""上传图片到抖音 IM 私信图床(VOD/zhenzhen 空间)。
成功返回 {uri(tos-cn-o-*), url, url_list, md5};失败返回 {"error": "..."}。
"""
del filename, content_type # VOD 按二进制上传,文件名仅用于本地存储
if not raw:
return {"error": "图片为空"}
try:
ak, sk, token, space = _fetch_im_upload_sts(session)
host, store_uri, jwt_auth, session_key = _vod_apply_upload_inner(
ak, sk, token, space, len(raw)
)
if not host or not store_uri or not jwt_auth:
return {"error": "申请上传地址失败:缺少 UploadHost/StoreUri/Auth"}
user_id = str(getattr(session, "my_uid", "") or "")
_vod_upload_binary(host, store_uri, jwt_auth, user_id, raw, session)
_vod_commit_upload_inner(ak, sk, token, space, session_key)
uri = store_uri.lstrip("/")
out: dict[str, Any] = {"uri": uri, "md5": hashlib.md5(raw).hexdigest()}
from .message_content import uri_to_cdn_urls
urls = uri_to_cdn_urls(uri)
if urls:
out["url_list"] = urls
out["url"] = urls[0]
logger.info("Uploaded IM image via VOD uri=%s host=%s space=%s", uri, host, space)
return out
except Exception as exc:
logger.warning("upload_im_image (VOD) failed: %s", exc)
return {"error": str(exc)}
def prepare_image_reply_spec(spec: dict[str, Any], session, upload_dir: str) -> tuple[dict[str, Any], str]:
"""若图片仍是本地地址,则上传到抖音 CDN 并补全 uri。返回 (spec, error)。"""
if spec.get("type") != "image":
return spec, ""
uri = str(spec.get("uri") or "").strip()
url = str(spec.get("url") or "").strip()
# 已经持有抖音 CDN 的 uri,说明图片早已上传完成(且经历过风控/转码),直接复用即可。
# 关键修复:不要因为 url 仍是本机预览地址(/api/media/...)而再次上传——
# 重复上传会拿到一个"刚提交、尚未完成风控/转码"的新 uri,发送时常被抖音以
# raw_check_code=1 / status_code=8003 拒绝(与是否互关无关)。
# 发送链路(build_msg_payload)只用 uri / url_list,从不使用这个本机 url,故本机 url 无害。
if uri and not is_local_media_url(uri):
return spec, ""
if url and is_douyin_cdn_url(url) and not uri:
return spec, ""
raw: bytes | None = None
filename = "image.jpg"
content_type = "image/jpeg"
# 尝试从各种本地 URL 形式读取文件
if is_local_media_url(url):
import os
# 先尝试 link-cards 目录(卡片封面图)
card_parsed = parse_link_card_media_path(url)
if card_parsed:
owner_id, fname = card_parsed
link_cards_dir = os.path.join(os.path.dirname(upload_dir), "link-cards")
path = os.path.join(link_cards_dir, str(owner_id), fname)
if os.path.isfile(path):
with open(path, "rb") as f:
raw = f.read()
filename = fname
if fname.lower().endswith(".png"):
content_type = "image/png"
elif fname.lower().endswith(".webp"):
content_type = "image/webp"
elif fname.lower().endswith(".gif"):
content_type = "image/gif"
# 如果文件不存在,尝试不带 .favicon 后缀
if raw is None and fname.endswith(".favicon.png"):
orig_path = os.path.join(link_cards_dir, str(owner_id), fname.replace(".favicon.png", ".png"))
if os.path.isfile(orig_path):
with open(orig_path, "rb") as f:
raw = f.read()
filename = fname.replace(".favicon.png", ".png")
content_type = "image/png"
# 再尝试 messages 目录
if raw is None:
parsed = parse_local_message_media_path(url)
if parsed:
_, fname = parsed
path = os.path.join(upload_dir, str(parsed[0]), fname)
if os.path.isfile(path):
with open(path, "rb") as f:
raw = f.read()
filename = fname
if fname.lower().endswith(".png"):
content_type = "image/png"
elif fname.lower().endswith(".webp"):
content_type = "image/webp"
elif fname.lower().endswith(".gif"):
content_type = "image/gif"
else:
# is_local_media_url 未识别,尝试强制从 URL 提取 link-cards 路径作为兜底
# (处理 KEFU_PUBLIC_BASE_URL 未配置或其他未知域名变体)
import os
card_parsed = parse_link_card_media_path(url)
if card_parsed:
owner_id, fname = card_parsed
link_cards_dir = os.path.join(os.path.dirname(upload_dir), "link-cards")
path = os.path.join(link_cards_dir, str(owner_id), fname)
if os.path.isfile(path):
with open(path, "rb") as f:
raw = f.read()
filename = fname
if fname.lower().endswith(".png"):
content_type = "image/png"
elif fname.lower().endswith(".webp"):
content_type = "image/webp"
elif fname.lower().endswith(".gif"):
content_type = "image/gif"
if raw is None:
if url.startswith("http") and not is_douyin_cdn_url(url):
return spec, "图片地址必须是抖音 CDN 或本地上传后的地址,外部 URL 无法用于 IM 发送"
if not url:
return spec, "缺少可上传的图片数据(URL 为空)"
return spec, f"缺少可上传的图片数据(URL: {url[:80]}"
uploaded = upload_im_image(session, raw, filename=filename, content_type=content_type)
if uploaded.get("error"):
return spec, uploaded["error"]
if not uploaded.get("uri"):
return spec, "抖音图片上传失败:未返回 uri"
merged = {
**spec,
**uploaded,
"type": "image",
"text": spec.get("text") or "[图片]",
}
return merged, ""
def upload_card_cover(cover_url: str, session, upload_dir: str) -> dict:
"""上传卡片封面图到抖音 CDN,返回 {uri, url, error}。"""
import os as _os
import re as _re
raw_url = (cover_url or "").strip()
if not raw_url:
return {"error": "封面 URL 为空"}
# 已经是抖音 CDN,直接复用
if is_douyin_cdn_url(raw_url):
return {"uri": "", "url": raw_url}
raw = None
filename = "cover.jpg"
content_type = "image/jpeg"
# 尝试从 link-cards 本地路径读取
card_parsed = parse_link_card_media_path(raw_url)
if card_parsed:
owner_id, fname = card_parsed
link_cards_dir = _os.path.join(_os.path.dirname(upload_dir), "link-cards")
path = _os.path.join(link_cards_dir, str(owner_id), fname)
if _os.path.isfile(path):
with open(path, "rb") as f:
raw = f.read()
filename = fname
if fname.lower().endswith(".png"):
content_type = "image/png"
elif fname.lower().endswith(".webp"):
content_type = "image/webp"
elif fname.endswith(".favicon.png"):
orig_path = _os.path.join(link_cards_dir, str(owner_id), fname.replace(".favicon.png", ".png"))
if _os.path.isfile(orig_path):
with open(orig_path, "rb") as f:
raw = f.read()
filename = fname.replace(".favicon.png", ".png")
content_type = "image/png"
# 尝试从 messages 目录读取
if raw is None:
msg_parsed = parse_local_message_media_path(raw_url)
if msg_parsed:
account_id, fname = msg_parsed
path = _os.path.join(upload_dir, str(account_id), fname)
if _os.path.isfile(path):
with open(path, "rb") as f:
raw = f.read()
filename = fname
if fname.lower().endswith(".png"):
content_type = "image/png"
elif fname.lower().endswith(".webp"):
content_type = "image/webp"
if raw is None:
return {"error": "无法从路径读取封面图: " + raw_url[:100]}
uploaded = upload_im_image(session, raw, filename=filename, content_type=content_type)
return uploaded
def _finalize_card_reply_spec(spec: dict[str, Any]) -> tuple[dict[str, Any], str]:
"""补全落地页 URL 并校验卡片是否可发送。"""
from .reply_payload import validate_card_send_spec
_fill_card_page_url(spec)
err = validate_card_send_spec(spec, require_cover_uri=True)
if err:
return spec, err
return spec, ""
def prepare_hyperlink_reply_spec(spec: dict[str, Any], session, upload_dir: str) -> tuple[dict[str, Any], str]:
"""发送超链接前:可选上传封面;校验直连 HTTPS URL(无落地页)。"""
if spec.get("type") != "hyperlink":
return spec, ""
from .reply_payload import validate_hyperlink_send_spec
merged = dict(spec)
cover_ref = str(
merged.get("cover_url") or merged.get("cover") or merged.get("image_path") or ""
).strip()
if not cover_ref:
return merged, "网页链接卡必须上传封面图"
existing_uri = str(merged.get("cover_uri") or merged.get("uri") or "").strip()
if existing_uri and existing_uri.startswith("tos-"):
merged["_cover_fresh_upload"] = False
err = validate_hyperlink_send_spec(merged, require_cover_uri=True)
return merged, err
cached = load_cached_cover_upload(cover_ref, upload_dir)
if cached:
_apply_cover_upload_to_spec(merged, cached)
merged["_cover_fresh_upload"] = False
err = validate_hyperlink_send_spec(merged, require_cover_uri=True)
return merged, err
uploaded = upload_card_cover(cover_ref, session, upload_dir)
if uploaded.get("error"):
return merged, str(uploaded["error"])
_apply_cover_upload_to_spec(merged, uploaded)
save_cached_cover_upload(cover_ref, upload_dir, uploaded)
merged["_cover_fresh_upload"] = True
err = validate_hyperlink_send_spec(merged, require_cover_uri=True)
return merged, err
def prepare_card_reply_spec(spec: dict[str, Any], session, upload_dir: str) -> tuple[dict[str, Any], str]:
"""发送链接卡片前:上传封面到抖音 CDN,并补全可公开访问的落地页 URL。
同一封面只上传一次并缓存 uri。重复上传会得到「刚提交、尚未过审」的新 tos uri,
发送时常被 raw_check_code=1 / status_code=8004 拒绝。
"""
if spec.get("type") != "card":
return spec, ""
merged = dict(spec)
cover_ref = str(
merged.get("cover_url")
or merged.get("cover")
or merged.get("image_path")
or ""
).strip()
if not cover_ref:
return merged, "缺少卡片封面图"
existing_uri = str(merged.get("cover_uri") or merged.get("uri") or "").strip()
if existing_uri and existing_uri.startswith("tos-"):
merged["_cover_fresh_upload"] = False
return _finalize_card_reply_spec(merged)
cached = load_cached_cover_upload(cover_ref, upload_dir)
if cached:
_apply_cover_upload_to_spec(merged, cached)
merged["_cover_fresh_upload"] = False
logger.info(
"Reusing cached card cover uri=%s key=%s",
str(cached.get("uri") or "")[:48],
normalize_cover_cache_key(cover_ref),
)
return _finalize_card_reply_spec(merged)
uploaded = upload_card_cover(cover_ref, session, upload_dir)
if uploaded.get("error"):
return merged, str(uploaded["error"])
_apply_cover_upload_to_spec(merged, uploaded)
save_cached_cover_upload(cover_ref, upload_dir, uploaded)
merged["_cover_fresh_upload"] = True
return _finalize_card_reply_spec(merged)