侧边栏在「监控」右边加了「运营」:账号列表 → 点进二级详情看该账号的数据。 【为什么是独立模块而不是监控的子视图】两者形状不同:监控是公开数据(点赞/收藏/评论/分享)的每轮快照+差分;运营是创作者后台按日期给出的曝光/观看/完播率/涨粉。凭据不同、采集方式也不同 —— 那边要浏览器登录态,这边是纯请求。硬塞进同一个模型会同时污染两边。 【扫码登录的关键差异】监控的扫码把登录态写进浏览器默认 profile(爬虫要复用)。运营要的是 cookie 字符串(纯请求够用),所以每次登录开一个**临时上下文**,扫完取出 cookie 就丢弃 —— 登第二个账号不会把第一个顶掉,也不影响监控那个登录态,十个账号互不干扰。 【决策依据】tools/probe_creator_api.py 的 Phase 0 实测:签名可自造(XYW_:MD5 → base64 → AES-128-CBC,与 xhshow 内置实现常量逐字节一致);主站 cookie 即可认证创作者后台;接口与参数已与真实页面对齐。 后端: - api/creator/models.py: creator_account / creator_note_stat。**复用 MonitorBase**,这样 create_all 与上一轮改成元数据驱动的 _ensure_columns 会自动覆盖新表 - api/creator/signing.py: XYW_ 签名,带三条实测结论(url= 前缀、appId=ugc、401 与 406 的区别) - api/creator/client.py: 纯 httpx 客户端。字段名尚未亲眼验证过,所以写成多别名匹配;解析不出来存 None 而非 0 - api/creator/service.py: 账号 CRUD 与同步。cookie 绝不进入对外结构,只给 has_cookie - api/creator/login.py: 临时上下文的扫码登录 - api/routers/creator.py: 8 条路由,全部带鉴权 前端: - 侧边栏「运营」+ OperationView(账号列表 → 二级详情)+ AddAccountDialog - 权限状态显眼呈现:pending 时照抄后台原话「已为您申请数据权限,次日可查看」,并说明此时同步返回 0 条是正常的,不是采集失败 测试:tests/test_creator_client.py 新增 48 例,含「cookie 不得出现在对外结构里」这条不变量,以及权限未生效时空壳响应的处理。
355 lines
13 KiB
Python
355 lines
13 KiB
Python
# -*- coding: utf-8 -*-
|
|
# Copyright (c) 2025 [email protected]
|
|
#
|
|
# This file is part of MediaCrawler project.
|
|
# Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/creator/service.py
|
|
# GitHub: https://github.com/NanmiCoder
|
|
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
|
|
#
|
|
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
|
|
# 1. 不得用于任何商业用途。
|
|
# 2. 使用时应遵守对应平台的使用条款和robots.txt规则。
|
|
# 3. 不得进行大规模爬取或对平台造成运营干扰。
|
|
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
|
|
# 5. 不得用于任何非法或不当的用途。
|
|
#
|
|
# 详细许可条款请参阅项目根目录下的LICENSE文件。
|
|
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
|
|
|
|
"""运营账号的增删查与数据同步。
|
|
|
|
一条贯穿全文件的规则:**cookie 是凭证,永远不出现在返回给上层的结构里。**
|
|
对外只给 `has_cookie` 这样的布尔量,与监控层对 cookie 的处理保持一致。
|
|
"""
|
|
|
|
import asyncio
|
|
from datetime import datetime, time, timedelta
|
|
from typing import Any, Dict, List, Optional
|
|
|
|
from sqlalchemy import delete, func, select
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from tools.time_util import get_current_timestamp
|
|
|
|
from .client import CreatorApiError, CreatorClient
|
|
from .models import (
|
|
ACCOUNT_ERROR,
|
|
ACCOUNT_EXPIRED,
|
|
ACCOUNT_OK,
|
|
PERMISSION_ACTIVE,
|
|
PERMISSION_MISSING,
|
|
PERMISSION_PENDING,
|
|
PERMISSION_UNKNOWN,
|
|
CreatorAccount,
|
|
CreatorNoteStat,
|
|
)
|
|
|
|
# 同步一次最多翻多少页。后台默认一页 10 条,200 页足以覆盖任何正常账号,
|
|
# 同时防止"接口不返回 has_more"时无限翻下去。
|
|
MAX_SYNC_PAGES = 200
|
|
PAGE_SIZE = 10
|
|
|
|
|
|
def _account_dict(account: CreatorAccount, note_count: int = 0) -> Dict[str, Any]:
|
|
"""账号的对外表示。**刻意不含 cookie。**"""
|
|
return {
|
|
"id": account.id,
|
|
"nickname": account.nickname,
|
|
"user_id": account.user_id,
|
|
"red_id": account.red_id,
|
|
"avatar": account.avatar,
|
|
"status": account.status,
|
|
"permission_status": account.permission_status,
|
|
# 后台原话照抄。"次日可查看"这类信息只能由它自己说,改写就失真了。
|
|
"permission_tip": account.permission_tip,
|
|
"last_checked_at": account.last_checked_at,
|
|
"last_synced_at": account.last_synced_at,
|
|
"last_error": account.last_error,
|
|
"has_cookie": bool(account.cookie),
|
|
"created_at": account.created_at,
|
|
"note_count": note_count,
|
|
}
|
|
|
|
|
|
async def list_accounts(session: AsyncSession) -> List[Dict[str, Any]]:
|
|
accounts = list(
|
|
(await session.scalars(select(CreatorAccount).order_by(CreatorAccount.id))).all()
|
|
)
|
|
counts = dict(
|
|
(
|
|
await session.execute(
|
|
select(CreatorNoteStat.account_id, func.count(func.distinct(CreatorNoteStat.note_id)))
|
|
.group_by(CreatorNoteStat.account_id)
|
|
)
|
|
).all()
|
|
)
|
|
return [_account_dict(account, counts.get(account.id, 0)) for account in accounts]
|
|
|
|
|
|
async def get_account(session: AsyncSession, account_id: int) -> CreatorAccount:
|
|
account = await session.get(CreatorAccount, account_id)
|
|
if account is None:
|
|
raise ValueError(f"账号 {account_id} 不存在")
|
|
return account
|
|
|
|
|
|
async def account_detail(session: AsyncSession, account_id: int) -> Dict[str, Any]:
|
|
account = await get_account(session, account_id)
|
|
notes = await latest_notes(session, account_id)
|
|
return {
|
|
"account": _account_dict(account, len(notes)),
|
|
"notes": notes,
|
|
"summary": _summarize(notes),
|
|
}
|
|
|
|
|
|
def _summarize(notes: List[Dict[str, Any]]) -> Dict[str, Any]:
|
|
"""账号级汇总。取最后一轮快照的累计值之和。"""
|
|
totals = {
|
|
key: 0
|
|
for key in ("exposure", "views", "likes", "comments", "favorites", "shares", "new_followers")
|
|
}
|
|
for note in notes:
|
|
for key in totals:
|
|
value = note.get(key)
|
|
if isinstance(value, (int, float)):
|
|
totals[key] += int(value)
|
|
return totals
|
|
|
|
|
|
async def latest_notes(session: AsyncSession, account_id: int) -> List[Dict[str, Any]]:
|
|
"""每个作品取**最近一次**快照。
|
|
|
|
表里保留全部历史(换个时点就是一条新行),但列表只该展示"现在",否则同一个
|
|
作品会在列表里出现多次。
|
|
"""
|
|
newest = (
|
|
select(
|
|
CreatorNoteStat.note_id,
|
|
func.max(CreatorNoteStat.captured_at).label("captured_at"),
|
|
)
|
|
.where(CreatorNoteStat.account_id == account_id)
|
|
.group_by(CreatorNoteStat.note_id)
|
|
.subquery()
|
|
)
|
|
rows = (
|
|
await session.scalars(
|
|
select(CreatorNoteStat)
|
|
.join(
|
|
newest,
|
|
(CreatorNoteStat.note_id == newest.c.note_id)
|
|
& (CreatorNoteStat.captured_at == newest.c.captured_at),
|
|
)
|
|
.where(CreatorNoteStat.account_id == account_id)
|
|
.order_by(CreatorNoteStat.publish_time.desc().nullslast())
|
|
)
|
|
).all()
|
|
return [_note_dict(row) for row in rows]
|
|
|
|
|
|
def _note_dict(row: CreatorNoteStat) -> Dict[str, Any]:
|
|
return {
|
|
"note_id": row.note_id,
|
|
"title": row.title,
|
|
"publish_time": row.publish_time,
|
|
"exposure": row.exposure,
|
|
"views": row.views,
|
|
"likes": row.likes,
|
|
"comments": row.comments,
|
|
"favorites": row.favorites,
|
|
"shares": row.shares,
|
|
"new_followers": row.new_followers,
|
|
"danmaku": row.danmaku,
|
|
"cover_ctr": row.cover_ctr,
|
|
"avg_watch_seconds": row.avg_watch_seconds,
|
|
"two_second_exit_rate": row.two_second_exit_rate,
|
|
"completion_rate": row.completion_rate,
|
|
"captured_at": row.captured_at,
|
|
}
|
|
|
|
|
|
async def upsert_account_from_cookie(session: AsyncSession, cookie: str) -> Dict[str, Any]:
|
|
"""用一份 cookie 识别并保存账号。
|
|
|
|
识别靠 `user/info` 而不是让用户填名字 —— 填错名字只会让后面所有数据对不上号。
|
|
已有同 `user_id` 的账号则更新它的 cookie(重新登录)。
|
|
"""
|
|
client = CreatorClient(cookie)
|
|
if not client.looks_authenticated:
|
|
raise ValueError("这份 cookie 里没有 a1,无法签名,请重新扫码")
|
|
|
|
try:
|
|
info = await client.fetch_user_info()
|
|
except CreatorApiError as exc:
|
|
raise ValueError(f"登录态无法使用:{exc}") from exc
|
|
|
|
if not info.get("user_id"):
|
|
raise ValueError("接口没有返回账号标识,可能登录态无效")
|
|
|
|
now = get_current_timestamp()
|
|
account = await session.scalar(
|
|
select(CreatorAccount).where(CreatorAccount.user_id == info["user_id"])
|
|
)
|
|
if account is None:
|
|
account = CreatorAccount(created_at=now)
|
|
session.add(account)
|
|
|
|
account.nickname = info.get("nickname") or account.nickname or "未命名账号"
|
|
account.user_id = info["user_id"]
|
|
account.red_id = info.get("red_id") or ""
|
|
account.avatar = info.get("avatar") or ""
|
|
account.cookie = cookie
|
|
account.status = ACCOUNT_OK
|
|
account.last_error = None
|
|
account.last_checked_at = now
|
|
account.updated_at = now
|
|
|
|
await session.flush()
|
|
|
|
# 顺手把权限状态也拉一次:新账号几乎必然处于"已申请、次日生效",
|
|
# 当场告诉用户,比让他明天再回来问要好。
|
|
await refresh_permission(session, account)
|
|
|
|
return _account_dict(account)
|
|
|
|
|
|
async def refresh_permission(session: AsyncSession, account: CreatorAccount) -> None:
|
|
"""查询并记录数据权限状态。失败不影响账号本身可用。"""
|
|
try:
|
|
permission = await CreatorClient(account.cookie).fetch_permission()
|
|
except CreatorApiError as exc:
|
|
if exc.status == 401:
|
|
account.status = ACCOUNT_EXPIRED
|
|
account.last_error = str(exc)
|
|
else:
|
|
account.last_error = str(exc)
|
|
account.updated_at = get_current_timestamp()
|
|
return
|
|
|
|
display = permission.get("display")
|
|
status = permission.get("status")
|
|
account.permission_tip = permission.get("tip") or ""
|
|
|
|
if display or status:
|
|
account.permission_status = PERMISSION_ACTIVE
|
|
elif account.permission_tip:
|
|
# 有提示语但未开通 —— 实测就是"已为您申请数据权限,次日可查看"。
|
|
account.permission_status = PERMISSION_PENDING
|
|
else:
|
|
account.permission_status = PERMISSION_MISSING
|
|
|
|
account.status = ACCOUNT_OK
|
|
account.last_error = None
|
|
account.last_checked_at = get_current_timestamp()
|
|
account.updated_at = account.last_checked_at
|
|
|
|
|
|
async def check_account(session: AsyncSession, account_id: int) -> Dict[str, Any]:
|
|
"""重新检测一个账号:登录态还在不在、权限开通没有。"""
|
|
account = await get_account(session, account_id)
|
|
await refresh_permission(session, account)
|
|
count = (
|
|
await session.execute(
|
|
select(func.count(func.distinct(CreatorNoteStat.note_id))).where(
|
|
CreatorNoteStat.account_id == account_id
|
|
)
|
|
)
|
|
).scalar() or 0
|
|
return _account_dict(account, count)
|
|
|
|
|
|
async def delete_account(session: AsyncSession, account_id: int) -> None:
|
|
account = await get_account(session, account_id)
|
|
await session.execute(
|
|
delete(CreatorNoteStat).where(CreatorNoteStat.account_id == account_id)
|
|
)
|
|
await session.delete(account)
|
|
|
|
|
|
def _day_bounds(days: int) -> tuple[int, int]:
|
|
"""最近 N 天的起止(毫秒)。与后台的按发布时间筛选对齐。"""
|
|
today = datetime.now()
|
|
end = int(datetime.combine(today.date(), time(23, 59, 59)).timestamp() * 1000)
|
|
start = int(
|
|
datetime.combine((today - timedelta(days=days)).date(), time(0, 0, 0)).timestamp() * 1000
|
|
)
|
|
return start, end
|
|
|
|
|
|
async def sync_account(
|
|
session: AsyncSession, account_id: int, days: int = 90
|
|
) -> Dict[str, Any]:
|
|
"""拉取一个账号的作品运营数据并落库。
|
|
|
|
权限未生效时接口返回的是**空壳成功**(`data.result` 里没有数据),不是错误 ——
|
|
所以"同步成功但 0 条"是正常结果,必须如实回报,不能让用户以为采集坏了。
|
|
"""
|
|
account = await get_account(session, account_id)
|
|
if not account.cookie:
|
|
raise ValueError("该账号没有可用的登录态,请重新扫码")
|
|
|
|
client = CreatorClient(account.cookie)
|
|
start_ms, end_ms = _day_bounds(days)
|
|
now = get_current_timestamp()
|
|
|
|
collected: List[Dict[str, Any]] = []
|
|
for page in range(1, MAX_SYNC_PAGES + 1):
|
|
try:
|
|
batch = await client.fetch_note_list(start_ms, end_ms, page_num=page, page_size=PAGE_SIZE)
|
|
except CreatorApiError as exc:
|
|
account.last_error = str(exc)
|
|
if exc.status == 401:
|
|
account.status = ACCOUNT_EXPIRED
|
|
account.updated_at = get_current_timestamp()
|
|
raise ValueError(f"同步失败:{exc}") from exc
|
|
|
|
collected.extend(batch)
|
|
if len(batch) < PAGE_SIZE:
|
|
break
|
|
await asyncio.sleep(0.6) # 对后台客气一点,这是自己的账号但仍是自动化访问
|
|
|
|
# 先删掉本时点可能存在的重复行,再写入 —— 表上有 (account, note, captured_at)
|
|
# 唯一索引,重复同步不该报错。
|
|
await session.execute(
|
|
delete(CreatorNoteStat).where(
|
|
CreatorNoteStat.account_id == account_id, CreatorNoteStat.captured_at == now
|
|
)
|
|
)
|
|
for note in collected:
|
|
if not note.get("note_id"):
|
|
continue
|
|
session.add(
|
|
CreatorNoteStat(
|
|
account_id=account_id,
|
|
note_id=note["note_id"],
|
|
title=note.get("title") or "",
|
|
publish_time=note.get("publish_time"),
|
|
exposure=note.get("exposure"),
|
|
views=note.get("views"),
|
|
likes=note.get("likes"),
|
|
comments=note.get("comments"),
|
|
favorites=note.get("favorites"),
|
|
shares=note.get("shares"),
|
|
new_followers=note.get("new_followers"),
|
|
danmaku=note.get("danmaku"),
|
|
cover_ctr=note.get("cover_ctr"),
|
|
avg_watch_seconds=note.get("avg_watch_seconds"),
|
|
two_second_exit_rate=note.get("two_second_exit_rate"),
|
|
completion_rate=note.get("completion_rate"),
|
|
captured_at=now,
|
|
)
|
|
)
|
|
|
|
account.last_synced_at = now
|
|
account.last_error = None
|
|
account.updated_at = now
|
|
await refresh_permission(session, account)
|
|
await session.flush()
|
|
|
|
return {
|
|
"account_id": account_id,
|
|
"fetched": len(collected),
|
|
"permission_status": account.permission_status,
|
|
"permission_tip": account.permission_tip,
|
|
}
|