common.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403
  1. #!/usr/bin/env python3
  2. """skill-video-lapian 共享库: 动态身份参数 + VOC网关 + S3 SigV4 上传.
  3. 铁律: 所有凭据/端点从环境或 fmode-identity.json 读取, 绝不硬编码。
  4. 每个数字生命用自己的 FEME_USERID / FEME_SESSION_TOKEN / FEME_NEWAPI_TOKEN。
  5. """
  6. import datetime
  7. import hashlib
  8. import hmac
  9. import json
  10. import os
  11. import re
  12. import ssl
  13. import sys
  14. import urllib.error
  15. import urllib.parse
  16. import urllib.request
  17. UA = ("Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 "
  18. "(KHTML, like Gecko) Chrome/126.0 Safari/537.36")
  19. _IDENTITY_CACHE = None
  20. # ---------------------------------------------------------------- identity --
  21. def _load_identity_file():
  22. """读取 /opt/data/fmode-identity.json (或 $FEME_IDENTITY_FILE)。"""
  23. global _IDENTITY_CACHE
  24. if _IDENTITY_CACHE is not None:
  25. return _IDENTITY_CACHE
  26. path = os.environ.get("FEME_IDENTITY_FILE", "/opt/data/fmode-identity.json")
  27. data = {}
  28. if os.path.exists(path):
  29. try:
  30. data = json.load(open(path, encoding="utf-8"))
  31. except Exception as e: # 文件坏了不致命, 继续用 env
  32. print(f"[common] WARN 解析 {path} 失败: {e}", file=sys.stderr)
  33. _IDENTITY_CACHE = data
  34. return data
  35. _TOKEN_SHAPES = {"newapi_token": ("sk-",), "session_token": ("r:",),
  36. "studio_url": ("http",), "userid": ("__ALNUM__",),
  37. "cloud_ak": ("__ALNUM__",), "cloud_sk": ("__ALNUM__",)}
  38. def _identity_get(*names, secret=False):
  39. """按优先级取身份参数: env(FEME_*/兼容旧名) > config.yaml > fmode-identity.json。
  40. 安全设计:
  41. - 值按形状校验(token 必须 sk-/r: 开头, userid 必须 objectId 形状) ——
  42. 部分宿主会对 env 值做脱敏或塞入无关值($PATH 等), 打码值/杂值一律当缺失回退
  43. - config.yaml 兜底: 雨飏容器实测(2026-08-28 采集脚本走它成功调 VOC 网关)
  44. """
  45. key = names[0]
  46. env_map = {
  47. "userid": ["FEME_USERID"],
  48. "session_token": ["FEME_SESSION_TOKEN", "FMODE_SESSION_TOKEN"],
  49. "newapi_token": ["FEME_NEWAPI_TOKEN", "FMODE_API_KEY", "FMODE_API_TOKEN"],
  50. "studio_url": ["FEME_STUDIO_URL"],
  51. "cloud_ak": ["CLOUD_SDK_AK"],
  52. "cloud_sk": ["CLOUD_SDK_SK"],
  53. }
  54. shapes = _TOKEN_SHAPES.get(key, ())
  55. ident = _load_identity_file()
  56. def _ok(v, shape):
  57. v = (v or "").strip()
  58. if not v:
  59. return None
  60. if shape == "__ALNUM__":
  61. # objectId/AK/SK 形状: 纯字母数字(挡 $PATH 等杂值与打码值)
  62. return v if re.fullmatch(r"[A-Za-z0-9]{4,64}", v) else None
  63. if shape and not v.startswith(shape):
  64. return None
  65. return v
  66. for n in env_map.get(key, [names]):
  67. v = _ok(os.environ.get(n), shapes[0] if shapes else ())
  68. if v:
  69. return v
  70. if key == "newapi_token":
  71. v = _ok(_token_from_config(), "sk-")
  72. if v:
  73. return v
  74. ident_key = {"session_token": "session_token", "userid": "userid",
  75. "newapi_token": "newapi_token"}.get(key)
  76. if ident_key:
  77. v = _ok(ident.get(ident_key), shapes[0] if shapes else ())
  78. if v:
  79. return v
  80. if key == "userid":
  81. raise IdentityError(
  82. "缺少身份参数 FEME_USERID。请 export FEME_USERID=<你的Parse objectId> "
  83. "或确保 /opt/data/fmode-identity.json 含 userid 字段 "
  84. "(每个生命必须用自己的 userid, 禁止借用他人空间)。")
  85. if key in ("cloud_ak", "cloud_sk"):
  86. n = "CLOUD_SDK_AK" if key == "cloud_ak" else "CLOUD_SDK_SK"
  87. v = (os.environ.get(n) or "").strip()
  88. if re.fullmatch(r"[A-Za-z0-9]{16,}", v or ""):
  89. return v
  90. return ""
  91. def _token_from_config():
  92. """从 /opt/data/config.yaml 读 api_key(含打码防护)。"""
  93. global _CONFIG_TOKEN
  94. if _CONFIG_TOKEN is not None:
  95. return _CONFIG_TOKEN
  96. _CONFIG_TOKEN = ""
  97. path = os.environ.get("FEME_CONFIG_FILE", "/opt/data/config.yaml")
  98. if not os.path.exists(path):
  99. return ""
  100. try:
  101. m = re.search(r'api_key:\s*["\']?(sk-[A-Za-z0-9_\-]+)', open(path, encoding="utf-8").read())
  102. if m:
  103. _CONFIG_TOKEN = m.group(1)
  104. except Exception as e:
  105. print(f"[common] WARN 读取 {path}: {e}", file=sys.stderr)
  106. return _CONFIG_TOKEN
  107. _CONFIG_TOKEN = None
  108. def _camel(snake):
  109. return "".join(p.capitalize() for p in snake.split("_"))
  110. class IdentityError(RuntimeError):
  111. pass
  112. def get_userid() -> str:
  113. uid = _identity_get("userid")
  114. if not uid:
  115. raise IdentityError(
  116. "缺少身份参数 FEME_USERID (或 fmode-identity.json.userid)。"
  117. "每个数字生命必须用自己的 userid, 禁止借用他人空间。")
  118. return uid
  119. def get_session_token() -> str:
  120. tok = _identity_get("session_token")
  121. if not tok:
  122. raise IdentityError("缺少 FEME_SESSION_TOKEN (飞马会话 r:...)")
  123. return tok
  124. def get_api_token() -> str:
  125. """API 计费 token (sk-..., 调 VOC 网关 / fmode-listen / 视觉模型都用它)。"""
  126. tok = _identity_get("newapi_token")
  127. if not tok:
  128. raise IdentityError("缺少 FEME_NEWAPI_TOKEN (sk-...)")
  129. return tok
  130. def get_studio_url() -> str:
  131. """自己的 studio 端点 (storage/credentials 上线后用)。"""
  132. url = _identity_get("studio_url")
  133. if url:
  134. return url.rstrip("/")
  135. uid = get_userid()
  136. # 雨飏001 端口约定; 其他生命容器请显式 export FEME_STUDIO_URL
  137. return "https://server.fmode.cn:19001" if uid == "sr2WiPsDyQ" else ""
  138. def get_storage_credentials(prefix: str = "") -> dict:
  139. """优先走 studio /api/storage/credentials (用户级临时凭证);
  140. 未上线(过渡期)回退容器级 CLOUD_SDK_AK/SK。返回 {access_key, secret_key}。"""
  141. studio = get_studio_url()
  142. if studio:
  143. try:
  144. body = json.dumps({"prefix": prefix}).encode()
  145. req = urllib.request.Request(
  146. f"{studio}/api/storage/credentials", data=body, method="POST",
  147. headers={"Content-Type": "application/json",
  148. "Authorization": f"Bearer {get_session_token()}"})
  149. resp = json.loads(urllib.request.urlopen(req, timeout=15).read())
  150. cred = resp.get("data") or resp
  151. if cred.get("access") and cred.get("secret"):
  152. return {"access_key": cred["access"], "secret_key": cred["secret"],
  153. "security_token": cred.get("securityToken") or "", "source": "studio"}
  154. except urllib.error.HTTPError as e:
  155. if e.code != 404:
  156. print(f"[common] WARN credentials API {e.code}, 回退容器 AK/SK", file=sys.stderr)
  157. except Exception as e:
  158. print(f"[common] WARN credentials API 不可用({e}), 回退容器 AK/SK", file=sys.stderr)
  159. ak, sk = _identity_get("cloud_ak"), _identity_get("cloud_sk")
  160. if ak and sk:
  161. return {"access_key": ak, "secret_key": sk, "security_token": "", "source": "container-env"}
  162. raise IdentityError("拿不到 S3 凭证: studio credentials 不可用且无 CLOUD_SDK_AK/SK")
  163. # --------------------------------------------------------------- VOC gateway --
  164. def voc_call(proxy_path: str, params=None, method="GET", retries=3):
  165. """调 https://server.fmode.cn/api/voc-social/<proxyPath>, Bearer sk-token。
  166. 解析层级注意: data.aweme_list (不是 data.data.aweme_list)。
  167. 返回 json 或 {"error": ...}。
  168. """
  169. base = os.environ.get("VOC_SOCIAL_GATEWAY", "https://server.fmode.cn/api/voc-social")
  170. url = f"{base}/{proxy_path.lstrip('/')}"
  171. if params and method.upper() == "GET":
  172. qs = urllib.parse.urlencode({k: v for k, v in params.items() if v is not None})
  173. if qs:
  174. url += "?" + qs
  175. headers = {"Authorization": f"Bearer {get_api_token()}",
  176. "Accept": "application/json", "User-Agent": UA}
  177. body = None
  178. if method.upper() == "POST":
  179. headers["Content-Type"] = "application/json"
  180. body = json.dumps(params or {}).encode()
  181. last = None
  182. for attempt in range(retries):
  183. req = urllib.request.Request(url, headers=headers, data=body,
  184. method=method.upper())
  185. try:
  186. with urllib.request.urlopen(req, timeout=90) as r:
  187. return json.loads(r.read().decode("utf-8", "replace"))
  188. except urllib.error.HTTPError as e:
  189. detail = e.read().decode("utf-8", "replace")[:300]
  190. last = f"HTTP {e.code}: {detail}"
  191. if e.code >= 500 and attempt < retries - 1:
  192. continue
  193. break
  194. except Exception as e:
  195. last = str(e)
  196. if attempt < retries - 1:
  197. continue
  198. break
  199. return {"error": last}
  200. def extract_aweme_detail(resp: dict) -> dict:
  201. """从 fetch_one_video 响应提取 aweme 对象(容多种层级)。"""
  202. data = resp.get("data") or {}
  203. aweme = data.get("aweme_detail")
  204. if not aweme:
  205. lst = data.get("aweme_list") or []
  206. aweme = lst[0] if lst else data
  207. return aweme or {}
  208. def aweme_summary(aweme: dict) -> dict:
  209. """提取报告需要的视频元数据。"""
  210. st = aweme.get("statistics") or {}
  211. author = aweme.get("author") or {}
  212. video = aweme.get("video") or {}
  213. play = (video.get("play_addr") or {}).get("url_list") or []
  214. download = (video.get("download_addr") or {}).get("url_list") or []
  215. tags = [t.get("hashtag_name") for t in (aweme.get("text_extra") or [])
  216. if t.get("hashtag_name")]
  217. ct = aweme.get("create_time")
  218. return {
  219. "aweme_id": aweme.get("aweme_id") or aweme.get("aweme_id_str") or "",
  220. "desc": aweme.get("desc") or "",
  221. "author": author.get("nickname") or "",
  222. "author_signature": author.get("signature") or "",
  223. "tags": tags,
  224. "duration_ms": aweme.get("duration") or video.get("duration") or 0,
  225. "create_time": datetime.datetime.fromtimestamp(ct).strftime("%Y-%m-%d %H:%M") if ct else "",
  226. "stats": {
  227. "digg": st.get("digg_count", 0), "comment": st.get("comment_count", 0),
  228. "collect": st.get("collect_count", 0), "share": st.get("share_count", 0),
  229. "play": st.get("play_count", 0),
  230. },
  231. "play_urls": play, "download_urls": download,
  232. }
  233. # ------------------------------------------------------------------- ffmpeg --
  234. def ffprobe_duration(path: str) -> float:
  235. import subprocess
  236. out = subprocess.run(
  237. ["ffprobe", "-v", "error", "-show_entries", "format=duration",
  238. "-of", "default=noprint_wrappers=1:nokey=1", str(path)],
  239. capture_output=True, text=True)
  240. try:
  241. return float(out.stdout.strip())
  242. except ValueError:
  243. return 0.0
  244. # ------------------------------------------------------------- S3 SigV4 PUT --
  245. # OBS(华为云) 必须 virtual-hosted addressing: bucket 放进 Host。
  246. S3_BUCKET = os.environ.get("S3_BUCKET", "storage-s3-nkkj")
  247. S3_REGION = os.environ.get("S3_REGION", "cn-north-4")
  248. S3_HOST = os.environ.get("S3_ENDPOINT", "https://obs.cn-north-4.myhuaweicloud.com") \
  249. .replace("https://", "").rstrip("/")
  250. def _sigv4(method: str, key: str, payload: bytes, cred: dict,
  251. content_type: str = "", query: str = "") -> dict:
  252. """生成 OBS S3 兼容 SigV4 头。cred = get_storage_credentials()。"""
  253. host = f"{S3_BUCKET}.{S3_HOST}"
  254. t = datetime.datetime.now(datetime.timezone.utc)
  255. amzdate = t.strftime("%Y%m%dT%H%M%SZ")
  256. datestamp = t.strftime("%Y%m%d")
  257. payload_hash = hashlib.sha256(payload).hexdigest()
  258. headers = {"host": host, "x-amz-content-sha256": payload_hash, "x-amz-date": amzdate}
  259. if content_type:
  260. headers["content-type"] = content_type
  261. if cred.get("security_token"):
  262. headers["x-amz-security-token"] = cred["security_token"]
  263. signed = ";".join(sorted(headers))
  264. canonical_headers = "".join(f"{k}:{headers[k]}\n" for k in sorted(headers))
  265. canonical = (f"{method}\n/{key}\n{query}\n{canonical_headers}\n{signed}\n{payload_hash}")
  266. scope = f"{datestamp}/{S3_REGION}/s3/aws4_request"
  267. sts = f"AWS4-HMAC-SHA256\n{amzdate}\n{scope}\n{hashlib.sha256(canonical.encode()).hexdigest()}"
  268. def hm(k, m):
  269. return hmac.new(k, m.encode(), hashlib.sha256).digest()
  270. k = hm(hm(hm(hm(("AWS4" + cred["secret_key"]).encode(), datestamp), S3_REGION), "s3"),
  271. "aws4_request")
  272. sig = hmac.new(k, sts.encode(), hashlib.sha256).hexdigest()
  273. auth = (f"AWS4-HMAC-SHA256 Credential={cred['access_key']}/{scope}, "
  274. f"SignedHeaders={signed}, Signature={sig}")
  275. out = {"Host": host, "Authorization": auth, "x-amz-date": amzdate,
  276. "x-amz-content-sha256": payload_hash, "x-amz-acl": "public-read"}
  277. if content_type:
  278. out["Content-Type"] = content_type
  279. if cred.get("security_token"):
  280. out["x-amz-security-token"] = cred["security_token"]
  281. return out
  282. def s3_put_bytes(key: str, payload: bytes, content_type: str = "application/octet-stream",
  283. public_read: bool = True) -> str:
  284. """上传字节到 S3_BUCKET, 返回可公开访问的 URL。"""
  285. cred = get_storage_credentials(f"user/{get_userid()}/")
  286. hdrs = _sigv4("PUT", key, payload, cred, content_type)
  287. if not public_read:
  288. hdrs.pop("x-amz-acl", None)
  289. url = f"https://{S3_BUCKET}.{S3_HOST}/{key}"
  290. req = urllib.request.Request(url, data=payload, method="PUT", headers=hdrs)
  291. with urllib.request.urlopen(req, timeout=300) as r:
  292. if r.status not in (200, 201):
  293. raise RuntimeError(f"S3 PUT {key} -> {r.status}")
  294. return public_url(key)
  295. def s3_put_file(path: str, key: str, content_type: str = None) -> str:
  296. ct = content_type or guess_content_type(path)
  297. with open(path, "rb") as f:
  298. return s3_put_bytes(key, f.read(), ct)
  299. def public_url(key: str) -> str:
  300. """公开访问 URL。s3.fmode.cn CNAME 生效前用 OBS 原生域名。"""
  301. if os.environ.get("S3_PUBLIC_BASE"):
  302. return f"{os.environ['S3_PUBLIC_BASE'].rstrip('/')}/{key}"
  303. return f"https://{S3_BUCKET}.{S3_HOST}/{key}"
  304. def guess_content_type(path: str) -> str:
  305. import mimetypes
  306. return mimetypes.guess_type(path)[0] or "application/octet-stream"
  307. # ----------------------------------------------------------------- helpers --
  308. def fmt_ts(ms: float) -> str:
  309. s = int(ms / 1000)
  310. return f"{s // 60}:{s % 60:02d}"
  311. def slugify_platform(name: str) -> str:
  312. """平台英文简称映射(douyin/xiaohongshu/weixin/bilibili/tiktok/shortdrama...)。"""
  313. m = {"抖音": "douyin", "douyin": "douyin", "tiktok": "tiktok",
  314. "小红书": "xiaohongshu", "xiaohongshu": "xiaohongshu", "rednote": "xiaohongshu",
  315. "微信": "weixin", "weixin": "weixin", "wechat": "weixin", "视频号": "weixin",
  316. "b站": "bilibili", "哔哩哔哩": "bilibili", "bilibili": "bilibili",
  317. "快手": "kuaishou", "kuaishou": "kuaishou", "短剧": "shortdrama",
  318. "shortdrama": "shortdrama"}
  319. return m.get((name or "").strip().lower(), "douyin")
  320. def download_file(url: str, dest: str, min_size: int = 100_000) -> int:
  321. """流式下载(完整, 不 Range 截断), 返回字节数。"""
  322. os.makedirs(os.path.dirname(dest) or ".", exist_ok=True)
  323. req = urllib.request.Request(url, headers={"User-Agent": UA})
  324. total = 0
  325. with urllib.request.urlopen(req, timeout=300) as r, open(dest, "wb") as f:
  326. while True:
  327. chunk = r.read(1 << 20)
  328. if not chunk:
  329. break
  330. f.write(chunk)
  331. total += len(chunk)
  332. if total < min_size:
  333. os.remove(dest)
  334. raise RuntimeError(f"下载不完整({total}B < {min_size}B), 已删除 {dest}")
  335. return total
  336. def e(text) -> str:
  337. """HTML 转义(模板里用)。"""
  338. from html import escape
  339. return escape(str(text or ""), quote=True)