| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427 |
- #!/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/<proxyPath>, 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。
- 规则(2026-08-30 用户确认, 硬性): 对外分发的链接 host 必须是 s3.fmode.cn,
- 不能用 OBS 原生域名(storage-s3-nkkj.obs.cn-north-4.myhuaweicloud.com) ——
- 报告里的相对路径资源只有经 s3.fmode.cn 反代才能被浏览器正常加载。
- S3_PUBLIC_BASE 环境变量仍可覆盖(高级用法)。签名/上传仍走 OBS 原生 host。
- """
- if os.environ.get("S3_PUBLIC_BASE"):
- return f"{os.environ['S3_PUBLIC_BASE'].rstrip('/')}/{key}"
- return f"https://s3.fmode.cn/{key}"
- def verify_public_host(key: str = "") -> bool:
- """启动自检: s3.fmode.cn 反代可达(200/302/403都算通, 404=桶无此对象但域名通)。
- 不可达时打印告警(链接规范仍不变, 交给人工处理 DNS)。"""
- url = f"https://s3.fmode.cn/{key or S3_BUCKET.replace('-s3','')}"
- try:
- req = urllib.request.Request(url, method="HEAD",
- headers={"User-Agent": "Mozilla/5.0"})
- with urllib.request.urlopen(req, timeout=12) as r:
- ok = r.status < 500
- except urllib.error.HTTPError as e:
- ok = e.code < 500
- except Exception as e:
- ok = False
- print(f"[common] WARN s3.fmode.cn 不可达({e}), 公开链接仍按规范生成, 请检查 DNS/反代",
- file=sys.stderr)
- return ok
- 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)
|