common.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338
  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. def _identity_get(*names, secret=False):
  36. """按优先级取身份参数: env(FEME_*) > env(兼容旧名) > fmode-identity.json。"""
  37. env_map = {
  38. "userid": ["FEME_USERID"],
  39. "session_token": ["FEME_SESSION_TOKEN", "FMODE_SESSION_TOKEN"],
  40. "newapi_token": ["FEME_NEWAPI_TOKEN", "FMODE_API_KEY", "FMODE_API_TOKEN"],
  41. "studio_url": ["FEME_STUDIO_URL"],
  42. "cloud_ak": ["CLOUD_SDK_AK"],
  43. "cloud_sk": ["CLOUD_SDK_SK"],
  44. }
  45. for env_names in env_map.get(names[0], [names]):
  46. for n in env_names:
  47. v = os.environ.get(n, "").strip()
  48. if v:
  49. return v
  50. ident = _load_identity_file()
  51. v = (ident.get(names[0]) or ident.get(_camel(names[0])) or "").strip()
  52. if v:
  53. return v
  54. return ""
  55. def _camel(snake):
  56. return "".join(p.capitalize() for p in snake.split("_"))
  57. class IdentityError(RuntimeError):
  58. pass
  59. def get_userid() -> str:
  60. uid = _identity_get("userid")
  61. if not uid:
  62. raise IdentityError(
  63. "缺少身份参数 FEME_USERID (或 fmode-identity.json.userid)。"
  64. "每个数字生命必须用自己的 userid, 禁止借用他人空间。")
  65. return uid
  66. def get_session_token() -> str:
  67. tok = _identity_get("session_token")
  68. if not tok:
  69. raise IdentityError("缺少 FEME_SESSION_TOKEN (飞马会话 r:...)")
  70. return tok
  71. def get_api_token() -> str:
  72. """API 计费 token (sk-..., 调 VOC 网关 / fmode-listen / 视觉模型都用它)。"""
  73. tok = _identity_get("newapi_token")
  74. if not tok:
  75. raise IdentityError("缺少 FEME_NEWAPI_TOKEN (sk-...)")
  76. return tok
  77. def get_studio_url() -> str:
  78. """自己的 studio 端点 (storage/credentials 上线后用)。"""
  79. url = _identity_get("studio_url")
  80. if url:
  81. return url.rstrip("/")
  82. uid = get_userid()
  83. # 雨飏001 端口约定; 其他生命容器请显式 export FEME_STUDIO_URL
  84. return "https://server.fmode.cn:19001" if uid == "sr2WiPsDyQ" else ""
  85. def get_storage_credentials(prefix: str = "") -> dict:
  86. """优先走 studio /api/storage/credentials (用户级临时凭证);
  87. 未上线(过渡期)回退容器级 CLOUD_SDK_AK/SK。返回 {access_key, secret_key}。"""
  88. studio = get_studio_url()
  89. if studio:
  90. try:
  91. body = json.dumps({"prefix": prefix}).encode()
  92. req = urllib.request.Request(
  93. f"{studio}/api/storage/credentials", data=body, method="POST",
  94. headers={"Content-Type": "application/json",
  95. "Authorization": f"Bearer {get_session_token()}"})
  96. resp = json.loads(urllib.request.urlopen(req, timeout=15).read())
  97. cred = resp.get("data") or resp
  98. if cred.get("access") and cred.get("secret"):
  99. return {"access_key": cred["access"], "secret_key": cred["secret"],
  100. "security_token": cred.get("securityToken") or "", "source": "studio"}
  101. except urllib.error.HTTPError as e:
  102. if e.code != 404:
  103. print(f"[common] WARN credentials API {e.code}, 回退容器 AK/SK", file=sys.stderr)
  104. except Exception as e:
  105. print(f"[common] WARN credentials API 不可用({e}), 回退容器 AK/SK", file=sys.stderr)
  106. ak, sk = _identity_get("cloud_ak"), _identity_get("cloud_sk")
  107. if ak and sk:
  108. return {"access_key": ak, "secret_key": sk, "security_token": "", "source": "container-env"}
  109. raise IdentityError("拿不到 S3 凭证: studio credentials 不可用且无 CLOUD_SDK_AK/SK")
  110. # --------------------------------------------------------------- VOC gateway --
  111. def voc_call(proxy_path: str, params=None, method="GET", retries=3):
  112. """调 https://server.fmode.cn/api/voc-social/<proxyPath>, Bearer sk-token。
  113. 解析层级注意: data.aweme_list (不是 data.data.aweme_list)。
  114. 返回 json 或 {"error": ...}。
  115. """
  116. base = os.environ.get("VOC_SOCIAL_GATEWAY", "https://server.fmode.cn/api/voc-social")
  117. url = f"{base}/{proxy_path.lstrip('/')}"
  118. if params and method.upper() == "GET":
  119. qs = urllib.parse.urlencode({k: v for k, v in params.items() if v is not None})
  120. if qs:
  121. url += "?" + qs
  122. headers = {"Authorization": f"Bearer {get_api_token()}",
  123. "Accept": "application/json", "User-Agent": UA}
  124. body = None
  125. if method.upper() == "POST":
  126. headers["Content-Type"] = "application/json"
  127. body = json.dumps(params or {}).encode()
  128. last = None
  129. for attempt in range(retries):
  130. req = urllib.request.Request(url, headers=headers, data=body,
  131. method=method.upper())
  132. try:
  133. with urllib.request.urlopen(req, timeout=90) as r:
  134. return json.loads(r.read().decode("utf-8", "replace"))
  135. except urllib.error.HTTPError as e:
  136. detail = e.read().decode("utf-8", "replace")[:300]
  137. last = f"HTTP {e.code}: {detail}"
  138. if e.code >= 500 and attempt < retries - 1:
  139. continue
  140. break
  141. except Exception as e:
  142. last = str(e)
  143. if attempt < retries - 1:
  144. continue
  145. break
  146. return {"error": last}
  147. def extract_aweme_detail(resp: dict) -> dict:
  148. """从 fetch_one_video 响应提取 aweme 对象(容多种层级)。"""
  149. data = resp.get("data") or {}
  150. aweme = data.get("aweme_detail")
  151. if not aweme:
  152. lst = data.get("aweme_list") or []
  153. aweme = lst[0] if lst else data
  154. return aweme or {}
  155. def aweme_summary(aweme: dict) -> dict:
  156. """提取报告需要的视频元数据。"""
  157. st = aweme.get("statistics") or {}
  158. author = aweme.get("author") or {}
  159. video = aweme.get("video") or {}
  160. play = (video.get("play_addr") or {}).get("url_list") or []
  161. download = (video.get("download_addr") or {}).get("url_list") or []
  162. tags = [t.get("hashtag_name") for t in (aweme.get("text_extra") or [])
  163. if t.get("hashtag_name")]
  164. ct = aweme.get("create_time")
  165. return {
  166. "aweme_id": aweme.get("aweme_id") or aweme.get("aweme_id_str") or "",
  167. "desc": aweme.get("desc") or "",
  168. "author": author.get("nickname") or "",
  169. "author_signature": author.get("signature") or "",
  170. "tags": tags,
  171. "duration_ms": aweme.get("duration") or video.get("duration") or 0,
  172. "create_time": datetime.datetime.fromtimestamp(ct).strftime("%Y-%m-%d %H:%M") if ct else "",
  173. "stats": {
  174. "digg": st.get("digg_count", 0), "comment": st.get("comment_count", 0),
  175. "collect": st.get("collect_count", 0), "share": st.get("share_count", 0),
  176. "play": st.get("play_count", 0),
  177. },
  178. "play_urls": play, "download_urls": download,
  179. }
  180. # ------------------------------------------------------------------- ffmpeg --
  181. def ffprobe_duration(path: str) -> float:
  182. import subprocess
  183. out = subprocess.run(
  184. ["ffprobe", "-v", "error", "-show_entries", "format=duration",
  185. "-of", "default=noprint_wrappers=1:nokey=1", str(path)],
  186. capture_output=True, text=True)
  187. try:
  188. return float(out.stdout.strip())
  189. except ValueError:
  190. return 0.0
  191. # ------------------------------------------------------------- S3 SigV4 PUT --
  192. # OBS(华为云) 必须 virtual-hosted addressing: bucket 放进 Host。
  193. S3_BUCKET = os.environ.get("S3_BUCKET", "storage-s3-nkkj")
  194. S3_REGION = os.environ.get("S3_REGION", "cn-north-4")
  195. S3_HOST = os.environ.get("S3_ENDPOINT", "https://obs.cn-north-4.myhuaweicloud.com") \
  196. .replace("https://", "").rstrip("/")
  197. def _sigv4(method: str, key: str, payload: bytes, cred: dict,
  198. content_type: str = "", query: str = "") -> dict:
  199. """生成 OBS S3 兼容 SigV4 头。cred = get_storage_credentials()。"""
  200. host = f"{S3_BUCKET}.{S3_HOST}"
  201. t = datetime.datetime.now(datetime.timezone.utc)
  202. amzdate = t.strftime("%Y%m%dT%H%M%SZ")
  203. datestamp = t.strftime("%Y%m%d")
  204. payload_hash = hashlib.sha256(payload).hexdigest()
  205. headers = {"host": host, "x-amz-content-sha256": payload_hash, "x-amz-date": amzdate}
  206. if content_type:
  207. headers["content-type"] = content_type
  208. if cred.get("security_token"):
  209. headers["x-amz-security-token"] = cred["security_token"]
  210. signed = ";".join(sorted(headers))
  211. canonical_headers = "".join(f"{k}:{headers[k]}\n" for k in sorted(headers))
  212. canonical = (f"{method}\n/{key}\n{query}\n{canonical_headers}\n{signed}\n{payload_hash}")
  213. scope = f"{datestamp}/{S3_REGION}/s3/aws4_request"
  214. sts = f"AWS4-HMAC-SHA256\n{amzdate}\n{scope}\n{hashlib.sha256(canonical.encode()).hexdigest()}"
  215. def hm(k, m):
  216. return hmac.new(k, m.encode(), hashlib.sha256).digest()
  217. k = hm(hm(hm(hm(("AWS4" + cred["secret_key"]).encode(), datestamp), S3_REGION), "s3"),
  218. "aws4_request")
  219. sig = hmac.new(k, sts.encode(), hashlib.sha256).hexdigest()
  220. auth = (f"AWS4-HMAC-SHA256 Credential={cred['access_key']}/{scope}, "
  221. f"SignedHeaders={signed}, Signature={sig}")
  222. out = {"Host": host, "Authorization": auth, "x-amz-date": amzdate,
  223. "x-amz-content-sha256": payload_hash, "x-amz-acl": "public-read"}
  224. if content_type:
  225. out["Content-Type"] = content_type
  226. if cred.get("security_token"):
  227. out["x-amz-security-token"] = cred["security_token"]
  228. return out
  229. def s3_put_bytes(key: str, payload: bytes, content_type: str = "application/octet-stream",
  230. public_read: bool = True) -> str:
  231. """上传字节到 S3_BUCKET, 返回可公开访问的 URL。"""
  232. cred = get_storage_credentials(f"user/{get_userid()}/")
  233. hdrs = _sigv4("PUT", key, payload, cred, content_type)
  234. if not public_read:
  235. hdrs.pop("x-amz-acl", None)
  236. url = f"https://{S3_BUCKET}.{S3_HOST}/{key}"
  237. req = urllib.request.Request(url, data=payload, method="PUT", headers=hdrs)
  238. with urllib.request.urlopen(req, timeout=300) as r:
  239. if r.status not in (200, 201):
  240. raise RuntimeError(f"S3 PUT {key} -> {r.status}")
  241. return public_url(key)
  242. def s3_put_file(path: str, key: str, content_type: str = None) -> str:
  243. ct = content_type or guess_content_type(path)
  244. with open(path, "rb") as f:
  245. return s3_put_bytes(key, f.read(), ct)
  246. def public_url(key: str) -> str:
  247. """公开访问 URL。s3.fmode.cn CNAME 生效前用 OBS 原生域名。"""
  248. if os.environ.get("S3_PUBLIC_BASE"):
  249. return f"{os.environ['S3_PUBLIC_BASE'].rstrip('/')}/{key}"
  250. return f"https://{S3_BUCKET}.{S3_HOST}/{key}"
  251. def guess_content_type(path: str) -> str:
  252. import mimetypes
  253. return mimetypes.guess_type(path)[0] or "application/octet-stream"
  254. # ----------------------------------------------------------------- helpers --
  255. def fmt_ts(ms: float) -> str:
  256. s = int(ms / 1000)
  257. return f"{s // 60}:{s % 60:02d}"
  258. def slugify_platform(name: str) -> str:
  259. """平台英文简称映射(douyin/xiaohongshu/weixin/bilibili/tiktok/shortdrama...)。"""
  260. m = {"抖音": "douyin", "douyin": "douyin", "tiktok": "tiktok",
  261. "小红书": "xiaohongshu", "xiaohongshu": "xiaohongshu", "rednote": "xiaohongshu",
  262. "微信": "weixin", "weixin": "weixin", "wechat": "weixin", "视频号": "weixin",
  263. "b站": "bilibili", "哔哩哔哩": "bilibili", "bilibili": "bilibili",
  264. "快手": "kuaishou", "kuaishou": "kuaishou", "短剧": "shortdrama",
  265. "shortdrama": "shortdrama"}
  266. return m.get((name or "").strip().lower(), "douyin")
  267. def download_file(url: str, dest: str, min_size: int = 100_000) -> int:
  268. """流式下载(完整, 不 Range 截断), 返回字节数。"""
  269. os.makedirs(os.path.dirname(dest) or ".", exist_ok=True)
  270. req = urllib.request.Request(url, headers={"User-Agent": UA})
  271. total = 0
  272. with urllib.request.urlopen(req, timeout=300) as r, open(dest, "wb") as f:
  273. while True:
  274. chunk = r.read(1 << 20)
  275. if not chunk:
  276. break
  277. f.write(chunk)
  278. total += len(chunk)
  279. if total < min_size:
  280. os.remove(dest)
  281. raise RuntimeError(f"下载不完整({total}B < {min_size}B), 已删除 {dest}")
  282. return total
  283. def e(text) -> str:
  284. """HTML 转义(模板里用)。"""
  285. from html import escape
  286. return escape(str(text or ""), quote=True)