#!/usr/bin/env python3 """skill-video-lapian 共享库: 动态身份参数 + VOC网关 + S3 SigV4 上传. 铁律: 所有凭据/端点从环境或 fmode-identity.json 读取, 绝不硬编码。 每个数字生命用自己的 FEME_USERID / FEME_SESSION_TOKEN / FEME_NEWAPI_TOKEN。 """ import datetime import hashlib import hmac import json import os import re import ssl import sys import urllib.error import urllib.parse import urllib.request UA = ("Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 " "(KHTML, like Gecko) Chrome/126.0 Safari/537.36") _IDENTITY_CACHE = None # ---------------------------------------------------------------- identity -- def _load_identity_file(): """读取 /opt/data/fmode-identity.json (或 $FEME_IDENTITY_FILE)。""" global _IDENTITY_CACHE if _IDENTITY_CACHE is not None: return _IDENTITY_CACHE path = os.environ.get("FEME_IDENTITY_FILE", "/opt/data/fmode-identity.json") data = {} if os.path.exists(path): try: data = json.load(open(path, encoding="utf-8")) except Exception as e: # 文件坏了不致命, 继续用 env print(f"[common] WARN 解析 {path} 失败: {e}", file=sys.stderr) _IDENTITY_CACHE = data return data _TOKEN_SHAPES = {"newapi_token": ("sk-",), "session_token": ("r:",), "studio_url": ("http",), "userid": ("__ALNUM__",), "cloud_ak": ("__ALNUM__",), "cloud_sk": ("__ALNUM__",)} def _identity_get(*names, secret=False): """按优先级取身份参数: env(FEME_*/兼容旧名) > config.yaml > fmode-identity.json。 安全设计: - 值按形状校验(token 必须 sk-/r: 开头, userid 必须 objectId 形状) —— 部分宿主会对 env 值做脱敏或塞入无关值($PATH 等), 打码值/杂值一律当缺失回退 - config.yaml 兜底: 雨飏容器实测(2026-08-28 采集脚本走它成功调 VOC 网关) """ key = names[0] env_map = { "userid": ["FEME_USERID"], "session_token": ["FEME_SESSION_TOKEN", "FMODE_SESSION_TOKEN"], "newapi_token": ["FEME_NEWAPI_TOKEN", "FMODE_API_KEY", "FMODE_API_TOKEN"], "studio_url": ["FEME_STUDIO_URL"], "cloud_ak": ["CLOUD_SDK_AK"], "cloud_sk": ["CLOUD_SDK_SK"], } shapes = _TOKEN_SHAPES.get(key, ()) ident = _load_identity_file() def _ok(v, shape): v = (v or "").strip() if not v: return None if shape == "__ALNUM__": # objectId/AK/SK 形状: 纯字母数字(挡 $PATH 等杂值与打码值) return v if re.fullmatch(r"[A-Za-z0-9]{4,64}", v) else None if shape and not v.startswith(shape): return None return v for n in env_map.get(key, [names]): v = _ok(os.environ.get(n), shapes[0] if shapes else ()) if v: return v if key == "newapi_token": v = _ok(_token_from_config(), "sk-") if v: return v ident_key = {"session_token": "session_token", "userid": "userid", "newapi_token": "newapi_token"}.get(key) if ident_key: v = _ok(ident.get(ident_key), shapes[0] if shapes else ()) if v: return v if key == "userid": raise IdentityError( "缺少身份参数 FEME_USERID。请 export FEME_USERID=<你的Parse objectId> " "或确保 /opt/data/fmode-identity.json 含 userid 字段 " "(每个生命必须用自己的 userid, 禁止借用他人空间)。") if key in ("cloud_ak", "cloud_sk"): n = "CLOUD_SDK_AK" if key == "cloud_ak" else "CLOUD_SDK_SK" v = (os.environ.get(n) or "").strip() if re.fullmatch(r"[A-Za-z0-9]{16,}", v or ""): return v return "" def _token_from_config(): """从 /opt/data/config.yaml 读 api_key(含打码防护)。""" global _CONFIG_TOKEN if _CONFIG_TOKEN is not None: return _CONFIG_TOKEN _CONFIG_TOKEN = "" path = os.environ.get("FEME_CONFIG_FILE", "/opt/data/config.yaml") if not os.path.exists(path): return "" try: m = re.search(r'api_key:\s*["\']?(sk-[A-Za-z0-9_\-]+)', open(path, encoding="utf-8").read()) if m: _CONFIG_TOKEN = m.group(1) except Exception as e: print(f"[common] WARN 读取 {path}: {e}", file=sys.stderr) return _CONFIG_TOKEN _CONFIG_TOKEN = None def _camel(snake): return "".join(p.capitalize() for p in snake.split("_")) class IdentityError(RuntimeError): pass def get_userid() -> str: uid = _identity_get("userid") if not uid: raise IdentityError( "缺少身份参数 FEME_USERID (或 fmode-identity.json.userid)。" "每个数字生命必须用自己的 userid, 禁止借用他人空间。") return uid def get_session_token() -> str: tok = _identity_get("session_token") if not tok: raise IdentityError("缺少 FEME_SESSION_TOKEN (飞马会话 r:...)") return tok def get_api_token() -> str: """API 计费 token (sk-..., 调 VOC 网关 / fmode-listen / 视觉模型都用它)。""" tok = _identity_get("newapi_token") if not tok: raise IdentityError("缺少 FEME_NEWAPI_TOKEN (sk-...)") return tok def get_studio_url() -> str: """自己的 studio 端点 (storage/credentials 上线后用)。""" url = _identity_get("studio_url") if url: return url.rstrip("/") uid = get_userid() # 雨飏001 端口约定; 其他生命容器请显式 export FEME_STUDIO_URL return "https://server.fmode.cn:19001" if uid == "sr2WiPsDyQ" else "" def get_storage_credentials(prefix: str = "") -> dict: """优先走 studio /api/storage/credentials (用户级临时凭证); 未上线(过渡期)回退容器级 CLOUD_SDK_AK/SK。返回 {access_key, secret_key}。""" studio = get_studio_url() if studio: try: body = json.dumps({"prefix": prefix}).encode() req = urllib.request.Request( f"{studio}/api/storage/credentials", data=body, method="POST", headers={"Content-Type": "application/json", "Authorization": f"Bearer {get_session_token()}"}) resp = json.loads(urllib.request.urlopen(req, timeout=15).read()) cred = resp.get("data") or resp if cred.get("access") and cred.get("secret"): return {"access_key": cred["access"], "secret_key": cred["secret"], "security_token": cred.get("securityToken") or "", "source": "studio"} except urllib.error.HTTPError as e: if e.code != 404: print(f"[common] WARN credentials API {e.code}, 回退容器 AK/SK", file=sys.stderr) except Exception as e: print(f"[common] WARN credentials API 不可用({e}), 回退容器 AK/SK", file=sys.stderr) ak, sk = _identity_get("cloud_ak"), _identity_get("cloud_sk") if ak and sk: return {"access_key": ak, "secret_key": sk, "security_token": "", "source": "container-env"} raise IdentityError("拿不到 S3 凭证: studio credentials 不可用且无 CLOUD_SDK_AK/SK") # --------------------------------------------------------------- VOC gateway -- def voc_call(proxy_path: str, params=None, method="GET", retries=3): """调 https://server.fmode.cn/api/voc-social/, Bearer sk-token。 解析层级注意: data.aweme_list (不是 data.data.aweme_list)。 返回 json 或 {"error": ...}。 """ base = os.environ.get("VOC_SOCIAL_GATEWAY", "https://server.fmode.cn/api/voc-social") url = f"{base}/{proxy_path.lstrip('/')}" if params and method.upper() == "GET": qs = urllib.parse.urlencode({k: v for k, v in params.items() if v is not None}) if qs: url += "?" + qs headers = {"Authorization": f"Bearer {get_api_token()}", "Accept": "application/json", "User-Agent": UA} body = None if method.upper() == "POST": headers["Content-Type"] = "application/json" body = json.dumps(params or {}).encode() last = None for attempt in range(retries): req = urllib.request.Request(url, headers=headers, data=body, method=method.upper()) try: with urllib.request.urlopen(req, timeout=90) as r: return json.loads(r.read().decode("utf-8", "replace")) except urllib.error.HTTPError as e: detail = e.read().decode("utf-8", "replace")[:300] last = f"HTTP {e.code}: {detail}" if e.code >= 500 and attempt < retries - 1: continue break except Exception as e: last = str(e) if attempt < retries - 1: continue break return {"error": last} def extract_aweme_detail(resp: dict) -> dict: """从 fetch_one_video 响应提取 aweme 对象(容多种层级)。""" data = resp.get("data") or {} aweme = data.get("aweme_detail") if not aweme: lst = data.get("aweme_list") or [] aweme = lst[0] if lst else data return aweme or {} def aweme_summary(aweme: dict) -> dict: """提取报告需要的视频元数据。""" st = aweme.get("statistics") or {} author = aweme.get("author") or {} video = aweme.get("video") or {} play = (video.get("play_addr") or {}).get("url_list") or [] download = (video.get("download_addr") or {}).get("url_list") or [] tags = [t.get("hashtag_name") for t in (aweme.get("text_extra") or []) if t.get("hashtag_name")] ct = aweme.get("create_time") return { "aweme_id": aweme.get("aweme_id") or aweme.get("aweme_id_str") or "", "desc": aweme.get("desc") or "", "author": author.get("nickname") or "", "author_signature": author.get("signature") or "", "tags": tags, "duration_ms": aweme.get("duration") or video.get("duration") or 0, "create_time": datetime.datetime.fromtimestamp(ct).strftime("%Y-%m-%d %H:%M") if ct else "", "stats": { "digg": st.get("digg_count", 0), "comment": st.get("comment_count", 0), "collect": st.get("collect_count", 0), "share": st.get("share_count", 0), "play": st.get("play_count", 0), }, "play_urls": play, "download_urls": download, } # ------------------------------------------------------------------- ffmpeg -- def ffprobe_duration(path: str) -> float: import subprocess out = subprocess.run( ["ffprobe", "-v", "error", "-show_entries", "format=duration", "-of", "default=noprint_wrappers=1:nokey=1", str(path)], capture_output=True, text=True) try: return float(out.stdout.strip()) except ValueError: return 0.0 # ------------------------------------------------------------- S3 SigV4 PUT -- # OBS(华为云) 必须 virtual-hosted addressing: bucket 放进 Host。 S3_BUCKET = os.environ.get("S3_BUCKET", "storage-s3-nkkj") S3_REGION = os.environ.get("S3_REGION", "cn-north-4") S3_HOST = os.environ.get("S3_ENDPOINT", "https://obs.cn-north-4.myhuaweicloud.com") \ .replace("https://", "").rstrip("/") def _sigv4(method: str, key: str, payload: bytes, cred: dict, content_type: str = "", query: str = "") -> dict: """生成 OBS S3 兼容 SigV4 头。cred = get_storage_credentials()。""" host = f"{S3_BUCKET}.{S3_HOST}" t = datetime.datetime.now(datetime.timezone.utc) amzdate = t.strftime("%Y%m%dT%H%M%SZ") datestamp = t.strftime("%Y%m%d") payload_hash = hashlib.sha256(payload).hexdigest() headers = {"host": host, "x-amz-content-sha256": payload_hash, "x-amz-date": amzdate} if content_type: headers["content-type"] = content_type if cred.get("security_token"): headers["x-amz-security-token"] = cred["security_token"] signed = ";".join(sorted(headers)) canonical_headers = "".join(f"{k}:{headers[k]}\n" for k in sorted(headers)) canonical = (f"{method}\n/{key}\n{query}\n{canonical_headers}\n{signed}\n{payload_hash}") scope = f"{datestamp}/{S3_REGION}/s3/aws4_request" sts = f"AWS4-HMAC-SHA256\n{amzdate}\n{scope}\n{hashlib.sha256(canonical.encode()).hexdigest()}" def hm(k, m): return hmac.new(k, m.encode(), hashlib.sha256).digest() k = hm(hm(hm(hm(("AWS4" + cred["secret_key"]).encode(), datestamp), S3_REGION), "s3"), "aws4_request") sig = hmac.new(k, sts.encode(), hashlib.sha256).hexdigest() auth = (f"AWS4-HMAC-SHA256 Credential={cred['access_key']}/{scope}, " f"SignedHeaders={signed}, Signature={sig}") out = {"Host": host, "Authorization": auth, "x-amz-date": amzdate, "x-amz-content-sha256": payload_hash, "x-amz-acl": "public-read"} if content_type: out["Content-Type"] = content_type if cred.get("security_token"): out["x-amz-security-token"] = cred["security_token"] return out def s3_put_bytes(key: str, payload: bytes, content_type: str = "application/octet-stream", public_read: bool = True) -> str: """上传字节到 S3_BUCKET, 返回可公开访问的 URL。""" cred = get_storage_credentials(f"user/{get_userid()}/") hdrs = _sigv4("PUT", key, payload, cred, content_type) if not public_read: hdrs.pop("x-amz-acl", None) url = f"https://{S3_BUCKET}.{S3_HOST}/{key}" req = urllib.request.Request(url, data=payload, method="PUT", headers=hdrs) with urllib.request.urlopen(req, timeout=300) as r: if r.status not in (200, 201): raise RuntimeError(f"S3 PUT {key} -> {r.status}") return public_url(key) def s3_put_file(path: str, key: str, content_type: str = None) -> str: ct = content_type or guess_content_type(path) with open(path, "rb") as f: return s3_put_bytes(key, f.read(), ct) def public_url(key: str) -> str: """公开访问 URL。s3.fmode.cn CNAME 生效前用 OBS 原生域名。""" if os.environ.get("S3_PUBLIC_BASE"): return f"{os.environ['S3_PUBLIC_BASE'].rstrip('/')}/{key}" return f"https://{S3_BUCKET}.{S3_HOST}/{key}" def guess_content_type(path: str) -> str: import mimetypes return mimetypes.guess_type(path)[0] or "application/octet-stream" # ----------------------------------------------------------------- helpers -- def fmt_ts(ms: float) -> str: s = int(ms / 1000) return f"{s // 60}:{s % 60:02d}" def slugify_platform(name: str) -> str: """平台英文简称映射(douyin/xiaohongshu/weixin/bilibili/tiktok/shortdrama...)。""" m = {"抖音": "douyin", "douyin": "douyin", "tiktok": "tiktok", "小红书": "xiaohongshu", "xiaohongshu": "xiaohongshu", "rednote": "xiaohongshu", "微信": "weixin", "weixin": "weixin", "wechat": "weixin", "视频号": "weixin", "b站": "bilibili", "哔哩哔哩": "bilibili", "bilibili": "bilibili", "快手": "kuaishou", "kuaishou": "kuaishou", "短剧": "shortdrama", "shortdrama": "shortdrama"} return m.get((name or "").strip().lower(), "douyin") def download_file(url: str, dest: str, min_size: int = 100_000) -> int: """流式下载(完整, 不 Range 截断), 返回字节数。""" os.makedirs(os.path.dirname(dest) or ".", exist_ok=True) req = urllib.request.Request(url, headers={"User-Agent": UA}) total = 0 with urllib.request.urlopen(req, timeout=300) as r, open(dest, "wb") as f: while True: chunk = r.read(1 << 20) if not chunk: break f.write(chunk) total += len(chunk) if total < min_size: os.remove(dest) raise RuntimeError(f"下载不完整({total}B < {min_size}B), 已删除 {dest}") return total def e(text) -> str: """HTML 转义(模板里用)。""" from html import escape return escape(str(text or ""), quote=True)