feat: persist history and work caches

This commit is contained in:
2026-09-07 16:40:32 +08:00
parent 07524a36cd
commit b42b2d33da
10 changed files with 761 additions and 201 deletions
+107 -83
View File
@@ -7,35 +7,10 @@ import time
from account_browser import BrowserManager, validate_config
from account_log import DailyLog, visible
from account_session import Session, SessionError
from account_store import Store, decode, notice_kind, notice_work_id
from account_store import Store, decode
from subscribe_notifications import notice_ids
def history_preview(notice):
kind = notice_kind(notice)
detail = notice.get(kind) or notice.get("favorite") or notice.get("collect") or {}
users = detail.get("from_user") or []
if isinstance(users, dict):
users = [users]
comment = detail.get("comment") or {}
if not users and comment.get("user"):
users = [comment["user"]]
aweme = detail.get("aweme") or notice.get("aweme") or {}
return {
"nid": notice.get("nid_str") or str(notice.get("nid", "")),
"create_time": notice.get("create_time") or notice.get("createTime") or "",
"kind": kind,
"work_id": notice_work_id(notice),
"work_desc": aweme.get("desc") or "",
"actor_uids": [str(user.get("uid", "")) for user in users],
"actor_names": [
user.get("nickname") or user.get("short_id") or "" for user in users
],
"comment": comment.get("text") or detail.get("content") or "",
"business": notice,
}
class Engine:
def __init__(self, root, playwright):
self.store = Store(root)
@@ -56,8 +31,6 @@ class Engine:
self.ready = set()
self.last_write = {}
self.verified = set()
self.history_cache = {}
self.works_cache = {}
self.logger = self.audit.logger
self.audit.record(
"历史摘要",
@@ -236,10 +209,39 @@ class Engine:
checked = heartbeat = time.monotonic()
pushed = added = 0
last_push = None
refreshed_at = self.store.setting(
"works_refreshed_at:" + ident
) or self.store.works_updated(ident)
next_works_refresh = (
refreshed_at
+ decode(self.store.account(ident)["rule"])[
"works_refresh_interval"
]
)
while ident in self.desired:
if time.monotonic() - checked >= 15:
await session.identity()
checked = time.monotonic()
refresh_interval = decode(self.store.account(ident)["rule"])[
"works_refresh_interval"
]
if refresh_interval and time.time() >= next_works_refresh:
try:
await self.fetch_works(ident, True, session)
except Exception as exc:
reason = (
str(exc)
if isinstance(exc, (ValueError, RuntimeError))
else type(exc).__name__
)
self.store.log(
"作品同步失败",
"定时增量刷新失败,本轮通知监听继续运行;到下个间隔再试",
logging.WARNING,
source=ident,
原因=reason,
)
next_works_refresh = time.time() + refresh_interval
waiting = await self.collect_details(ident, session)
self.state(
ident,
@@ -418,22 +420,24 @@ class Engine:
if account["role"] != "main" or account["uid"] is None:
raise ValueError("请选择已绑定身份的大号")
session = await self.open(ident)
known = self.store.cached_notice_ids(ident)
self.store.log(
"历史获取", "开始只读分页获取全部互动通知;不标记已读", source=ident
)
notices = await session.history_notices()
self.history_cache[ident] = {
notice.get("nid_str") or str(notice.get("nid", "")): notice
for notice in notices
}
previews = [history_preview(notice) for notice in notices]
self.store.log(
"历史获取",
"全部可用历史事件已加载,等待用户预览和勾选;尚未创建任务",
"历史同步",
"开始增量同步历史事件;首次使用会全量缓存,不标记已读",
source=ident,
事件数=len(previews),
本地已有=len(known),
)
return {"account": ident, "items": previews}
notices = await session.history_notices(known)
self.store.cache_notices(ident, notices, "history")
items = self.store.cached_notices(ident)
self.store.log(
"历史同步",
"历史事件已保存到本地,列表已按当前作品范围过滤;尚未创建任务",
source=ident,
本次返回=len(notices),
当前可选=len(items),
)
return {"account": ident, "items": items}
def enqueue_history(self, ident, ids):
if (
@@ -446,18 +450,19 @@ class Engine:
):
raise ValueError("历史事件选择无效")
ids = list(dict.fromkeys(ids))
cache = self.history_cache.get(ident)
if cache is None:
raise ValueError("历史事件预览已失效,请重新获取")
if any(item not in cache for item in ids):
raise ValueError("选择包含未加载的历史事件")
available = {item["nid"] for item in self.store.cached_notices(ident)}
if any(item not in available for item in ids):
raise ValueError("选择包含未缓存或不属于当前作品范围的历史事件")
rule = decode(self.store.account(ident)["rule"])
if not rule["enabled"] or not (rule["follow"] or rule["dm"]):
raise ValueError("请先启用大号规则并至少选择关注或私信")
before = self.store.db.execute("SELECT count(*) FROM tasks").fetchone()[0]
added = sum(
self.store.ingest(ident, cache[item], origin="history") for item in ids
)
added = 0
for nid in ids:
notice = self.store.cached_notice(ident, nid)
if notice is None:
raise ValueError("历史事件缓存已损坏,请重新同步")
added += self.store.ingest(ident, notice, origin="history")
tasks = (
self.store.db.execute("SELECT count(*) FROM tasks").fetchone()[0] - before
)
@@ -471,51 +476,65 @@ class Engine:
)
return {"selected": len(ids), "events": added, "tasks": tasks}
async def fetch_works(self, ident, all_pages):
async def fetch_works(self, ident, force=False, session=None):
account = self.store.account(ident)
if account["role"] != "main" or account["uid"] is None:
raise ValueError("请选择已绑定身份的大号")
if type(all_pages) is not bool:
raise ValueError("作品获取范围无效")
session = await self.open(ident)
self.store.log(
"作品获取",
"开始只读获取全部历史作品" if all_pages else "开始只读获取近期作品",
source=ident,
)
works = await session.works(all_pages)
self.works_cache[ident] = {work["aweme_id"]: work for work in works}
self.store.log(
"作品获取",
"作品列表已加载,等待用户选择监控范围",
source=ident,
作品数=len(works),
范围="全部历史" if all_pages else "近期",
)
if type(force) is not bool:
raise ValueError("作品刷新参数无效")
known = self.store.cached_work_ids(ident)
refreshed = 0
full = not known
if full or force:
session = session or await self.open(ident)
self.store.log(
"作品同步",
"首次建立全部历史作品缓存" if full else "开始增量刷新作品缓存",
source=ident,
本地已有=len(known),
)
works = await session.works(full, known if known else None)
refreshed = self.store.cache_works(ident, works)
self.store.set_setting("works_refreshed_at:" + ident, time.time())
items = self.store.cached_works(ident)
rule = decode(account["rule"])
self.store.log(
"作品同步",
"作品缓存可用,等待用户选择监控范围",
source=ident,
本次更新=refreshed,
本地总数=len(items),
)
return {
"account": ident,
"items": works,
"items": items,
"mode": rule["work_mode"],
"selected": rule["work_ids"],
"all_pages": all_pages,
"refresh_interval": rule["works_refresh_interval"],
"refreshed": refreshed,
}
def save_work_filter(self, ident, mode, ids):
def save_work_filter(self, ident, mode, ids, refresh_interval):
if mode not in ("all", "selected") or not isinstance(ids, list):
raise ValueError("作品监控选择无效")
ids = list(dict.fromkeys(ids))
if mode == "selected":
if not ids:
raise ValueError("指定作品模式至少选择一个作品")
cache = self.works_cache.get(ident)
existing = decode(self.store.account(ident)["rule"])["work_ids"]
if cache is None or any(
item not in cache and item not in existing for item in ids
):
raise ValueError("作品列表已失效,请重新获取")
self.store.set_work_filter(ident, mode, ids if mode == "selected" else [])
return {"mode": mode, "count": len(ids) if mode == "selected" else 0}
cached = self.store.cached_work_ids(ident)
if any(item not in cached for item in ids):
raise ValueError("选择包含未缓存的作品,请强制刷新后重试")
self.store.set_work_filter(
ident,
mode,
ids if mode == "selected" else [],
refresh_interval,
)
return {
"mode": mode,
"count": len(ids) if mode == "selected" else 0,
"refresh_interval": refresh_interval,
}
async def command(self, name, data):
labels = {
@@ -533,8 +552,8 @@ class Engine:
"delete_worker": "删除小号",
"history_fetch": "获取全部历史事件",
"history_enqueue": "确认历史事件操作",
"works_recent": "获取近期作品",
"works_all": "获取全部历史作品",
"works_open": "打开作品监控",
"works_refresh": "强制增量刷新作品",
"work_filter": "保存作品监控范围",
}
label = labels.get(name, "未知操作")
@@ -578,10 +597,15 @@ class Engine:
return await self.fetch_history(data["id"])
elif name == "history_enqueue":
return self.enqueue_history(data["id"], data.get("ids"))
elif name in ("works_recent", "works_all"):
return await self.fetch_works(data["id"], name == "works_all")
elif name in ("works_open", "works_refresh"):
return await self.fetch_works(data["id"], name == "works_refresh")
elif name == "work_filter":
return self.save_work_filter(data["id"], data.get("mode"), data.get("ids"))
return self.save_work_filter(
data["id"],
data.get("mode"),
data.get("ids"),
data.get("refresh_interval"),
)
elif name == "move":
self.store.move(data["id"], data.get("owner"))
owner = data.get("owner")
+17 -9
View File
@@ -232,8 +232,9 @@ class Session:
raise SessionError("通知详情返回了未请求的 ID,已停止处理")
return notices
async def history_notices(self):
async def history_notices(self, stop_ids=None):
notices = {}
stop_ids = set(stop_ids or ())
min_time = max_time = 0
seen_cursors = set()
for _ in range(1000):
@@ -256,13 +257,14 @@ class Session:
not isinstance(row, dict) for row in rows
):
raise SessionError("历史通知列表格式无效")
reached_cache = False
for row in rows:
if str(row.get("user_id")) != self.uid:
raise SessionError("历史通知所属身份不符")
nid = row.get("nid_str") or row.get("nid")
validate_uid(nid)
nid = validate_uid(row.get("nid_str") or row.get("nid"))
reached_cache = reached_cache or nid in stop_ids
notices.setdefault(nid, row)
if not payload.get("has_more"):
if reached_cache or not payload.get("has_more"):
return list(notices.values())
cursor = (payload.get("min_time"), payload.get("max_time"))
if cursor in seen_cursors or cursor == (min_time, max_time):
@@ -273,15 +275,16 @@ class Session:
min_time, max_time = cursor
raise SessionError("历史通知超过 50000 条,已停止以避免无限分页")
async def works(self, all_pages=False):
async def works(self, all_pages=False, stop_ids=None):
profile = await self.identity()
stop_ids = set(stop_ids or ())
sec_uid = profile.get("sec_uid") or profile.get("secUid")
if not isinstance(sec_uid, str) or not sec_uid:
raise SessionError("当前账号缺少作品列表身份参数")
works = {}
cursor = 0
seen_cursors = set()
pages = 1000 if all_pages else 1
pages = 1000 if all_pages or stop_ids else 1
for _ in range(pages):
response = await self.json(
works_request_script(sec_uid, cursor), main_world=True
@@ -302,12 +305,13 @@ class Session:
not isinstance(row, dict) for row in rows
):
raise SessionError("作品列表格式无效")
reached_cache = False
for row in rows:
author = row.get("author") or {}
if str(author.get("uid")) != self.uid:
raise SessionError("作品列表所属身份不符")
ident = row.get("aweme_id") or row.get("awemeId")
validate_uid(ident)
ident = validate_uid(row.get("aweme_id") or row.get("awemeId"))
reached_cache = reached_cache or ident in stop_ids
cover = ((row.get("video") or {}).get("cover") or {}).get(
"url_list"
) or []
@@ -324,7 +328,11 @@ class Session:
"business": row,
},
)
if not all_pages or not payload.get("has_more"):
if (
reached_cache
or (not all_pages and not stop_ids)
or not payload.get("has_more")
):
return list(works.values())
next_cursor = payload.get("max_cursor")
if type(next_cursor) not in (int, float) or next_cursor < 0:
+207 -4
View File
@@ -18,6 +18,7 @@ def validate_rule(rule):
rule.setdefault("cooldown", 14400)
rule.setdefault("work_mode", "all")
rule.setdefault("work_ids", [])
rule.setdefault("works_refresh_interval", 3600)
kinds = rule.get("kinds", [])
if not isinstance(kinds, list) or any(
k not in ("digg", "follow", "comment", "general_notice") for k in kinds
@@ -48,6 +49,11 @@ def validate_rule(rule):
):
raise ValueError("作品 ID 列表无效")
rule["work_ids"] = list(dict.fromkeys(rule["work_ids"]))
if (
type(rule["works_refresh_interval"]) is not int
or not 0 <= rule["works_refresh_interval"] <= 604800
):
raise ValueError("作品自动刷新间隔需为 0..604800 秒")
return rule
@@ -62,6 +68,7 @@ DEFAULT_RULE = {
"cooldown": 14400,
"work_mode": "all",
"work_ids": [],
"works_refresh_interval": 3600,
}
@@ -102,6 +109,31 @@ def notice_kind(notice):
return "general_notice"
def notice_preview(notice):
kind = notice_kind(notice)
detail = notice.get(kind) or notice.get("favorite") or notice.get("collect") or {}
users = detail.get("from_user") or []
if isinstance(users, dict):
users = [users]
comment = detail.get("comment") or {}
if not users and comment.get("user"):
users = [comment["user"]]
aweme = detail.get("aweme") or notice.get("aweme") or {}
return {
"nid": notice.get("nid_str") or str(notice.get("nid", "")),
"create_time": notice.get("create_time") or notice.get("createTime") or "",
"kind": kind,
"work_id": notice_work_id(notice),
"work_desc": aweme.get("desc") or "",
"actor_uids": [str(user.get("uid", "")) for user in users],
"actor_names": [
user.get("nickname") or user.get("short_id") or "" for user in users
],
"comment": comment.get("text") or detail.get("content") or "",
"business": notice,
}
class Store:
def __init__(self, root, audit=None):
self.audit = audit
@@ -147,6 +179,23 @@ class Store:
source TEXT NOT NULL REFERENCES accounts(id), target TEXT NOT NULL,
last_at REAL NOT NULL, event INTEGER NOT NULL REFERENCES events(id),
PRIMARY KEY(source,target));
CREATE TABLE IF NOT EXISTS cached_works (
source TEXT NOT NULL REFERENCES accounts(id), aweme_id TEXT NOT NULL,
create_time REAL NOT NULL DEFAULT 0, description TEXT NOT NULL DEFAULT '',
statistics TEXT NOT NULL DEFAULT '{}', cover TEXT NOT NULL DEFAULT '',
business TEXT NOT NULL, updated REAL NOT NULL,
PRIMARY KEY(source,aweme_id));
CREATE INDEX IF NOT EXISTS cached_works_time ON cached_works(source,create_time DESC);
CREATE TABLE IF NOT EXISTS cached_notices (
source TEXT NOT NULL REFERENCES accounts(id), nid TEXT NOT NULL,
create_time REAL NOT NULL DEFAULT 0, kind TEXT NOT NULL,
work_id TEXT NOT NULL DEFAULT '', actor_uids TEXT NOT NULL DEFAULT '[]',
actor_names TEXT NOT NULL DEFAULT '[]', comment TEXT NOT NULL DEFAULT '',
work_desc TEXT NOT NULL DEFAULT '', business TEXT NOT NULL,
origin TEXT NOT NULL CHECK(origin IN ('live','history')), updated REAL NOT NULL,
PRIMARY KEY(source,nid));
CREATE INDEX IF NOT EXISTS cached_notices_time ON cached_notices(source,create_time DESC);
CREATE INDEX IF NOT EXISTS cached_notices_work ON cached_notices(source,work_id,kind,create_time DESC);
DROP INDEX IF EXISTS worker_single_running;
""")
self.migrate_accounts()
@@ -351,7 +400,13 @@ class Store:
"首次身份核验通过,绑定已锁定;业务资料按原文输出,认证凭据保持隐藏",
account=ident,
UID=uid,
平台资料=saved,
昵称=saved.get("nickname") or "",
抖音号=saved.get("douyin_id") or "",
头像URL=saved.get("avatar_url") or "",
关注数=saved.get("following_count") or 0,
粉丝数=saved.get("follower_count") or 0,
获赞数=saved.get("total_favorited") or 0,
作品数=saved.get("aweme_count") or 0,
)
self.db.execute(
"UPDATE accounts SET uid=?,nickname=?,avatar=?,profile=? WHERE id=?",
@@ -383,16 +438,164 @@ class Store:
动作间隔秒=rule["interval"],
同组UID冷却秒=rule["cooldown"],
作品监控模式=rule["work_mode"],
作品ID=rule["work_ids"],
作品ID=",".join(rule["work_ids"]),
私信正文=rule["text"],
)
def set_work_filter(self, ident, mode, ids):
def set_work_filter(self, ident, mode, ids, refresh_interval=None):
rule = decode(self.account(ident)["rule"])
rule["work_mode"] = mode
rule["work_ids"] = ids
if refresh_interval is not None:
rule["works_refresh_interval"] = refresh_interval
self.set_rule(ident, rule)
def cache_works(self, source, works):
self.account(source)
now = time.time()
with self.transaction():
for work in works:
ident = validate_uid(work.get("aweme_id"))
created = work.get("create_time") or 0
if type(created) not in (int, float):
created = 0
self.db.execute(
"""INSERT INTO cached_works(source,aweme_id,create_time,description,statistics,cover,business,updated)
VALUES (?,?,?,?,?,?,?,?) ON CONFLICT(source,aweme_id) DO UPDATE SET
create_time=excluded.create_time,description=excluded.description,
statistics=excluded.statistics,cover=excluded.cover,business=excluded.business,updated=excluded.updated""",
(
source,
ident,
created,
str(work.get("desc") or ""),
json.dumps(work.get("statistics") or {}, ensure_ascii=False),
str(work.get("cover") or ""),
json.dumps(
visible(work.get("business") or work), ensure_ascii=False
),
now,
),
)
return len(works)
def cached_work_ids(self, source):
return {
row["aweme_id"]
for row in self.db.execute(
"SELECT aweme_id FROM cached_works WHERE source=?", (source,)
)
}
def cached_works(self, source):
self.account(source)
return [
{
"aweme_id": row["aweme_id"],
"desc": row["description"],
"create_time": row["create_time"],
"statistics": decode(row["statistics"]),
"cover": row["cover"],
"updated": row["updated"],
}
for row in self.db.execute(
"SELECT * FROM cached_works WHERE source=? ORDER BY create_time DESC,aweme_id DESC",
(source,),
)
]
def works_updated(self, source):
row = self.db.execute(
"SELECT max(updated) AS updated FROM cached_works WHERE source=?", (source,)
).fetchone()
return row["updated"] or 0
def cache_notices(self, source, notices, origin):
account = self.account(source)
if origin not in ("live", "history"):
raise ValueError("通知缓存来源无效")
now = time.time()
with self.transaction():
for notice in notices:
if str(notice.get("user_id")) != account["uid"]:
raise ValueError("通知缓存所属身份不符")
preview = notice_preview(notice)
nid = validate_uid(preview["nid"])
created = preview["create_time"] or 0
if type(created) not in (int, float):
created = 0
if created > 10_000_000_000:
created /= 1000
self.db.execute(
"""INSERT INTO cached_notices(source,nid,create_time,kind,work_id,actor_uids,actor_names,comment,work_desc,business,origin,updated)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(source,nid) DO UPDATE SET
create_time=excluded.create_time,kind=excluded.kind,work_id=excluded.work_id,
actor_uids=excluded.actor_uids,actor_names=excluded.actor_names,
comment=excluded.comment,work_desc=excluded.work_desc,business=excluded.business,
origin=CASE WHEN cached_notices.origin='live' OR excluded.origin='live' THEN 'live' ELSE 'history' END,
updated=excluded.updated""",
(
source,
nid,
created,
preview["kind"],
preview["work_id"],
json.dumps(preview["actor_uids"], ensure_ascii=False),
json.dumps(preview["actor_names"], ensure_ascii=False),
preview["comment"],
preview["work_desc"],
json.dumps(visible(notice), ensure_ascii=False),
origin,
now,
),
)
return len(notices)
def cached_notice_ids(self, source):
return {
row["nid"]
for row in self.db.execute(
"SELECT nid FROM cached_notices WHERE source=?", (source,)
)
}
def cached_notices(self, source):
rule = decode(self.account(source)["rule"])
rows = self.db.execute(
"SELECT * FROM cached_notices WHERE source=? ORDER BY create_time DESC,nid DESC",
(source,),
)
selected = set(rule["work_ids"])
result = []
for row in rows:
if (
rule["work_mode"] == "selected"
and row["kind"] != "follow"
and row["work_id"] not in selected
):
continue
result.append(
{
"nid": row["nid"],
"create_time": row["create_time"],
"kind": row["kind"],
"work_id": row["work_id"],
"work_desc": row["work_desc"],
"actor_uids": decode(row["actor_uids"]),
"actor_names": decode(row["actor_names"]),
"comment": row["comment"],
"origin": row["origin"],
}
)
return result
def cached_notice(self, source, nid):
row = self.db.execute(
"SELECT business FROM cached_notices WHERE source=? AND nid=?",
(source, nid),
).fetchone()
return decode(row["business"]) if row else None
def move(self, worker, owner):
with self.transaction():
account = self.account(worker)
@@ -528,6 +731,7 @@ class Store:
if account["role"] != "main" or str(notice.get("user_id")) != account["uid"]:
raise ValueError("通知所属身份不符")
nid = validate_uid(notice.get("nid_str") or str(notice.get("nid", "")))
self.cache_notices(source, [notice], origin)
kind = notice_kind(notice)
detail = (
notice.get(kind) or notice.get("favorite") or notice.get("collect") or {}
@@ -612,7 +816,6 @@ class Store:
}[kind],
作品ID=work_id,
目标数量=len(targets),
业务数据=notice,
)
self.log(
"规则匹配" if enabled else "通知跳过",
+178 -71
View File
@@ -6,6 +6,7 @@ import json
import queue
import sys
import tempfile
from datetime import datetime
from pathlib import Path
from patchright.async_api import async_playwright
@@ -44,6 +45,55 @@ from account_log import UI_LIMIT, DailyLog
from account_store import decode, validate_rule
def date_text(value):
try:
stamp = float(value)
except (TypeError, ValueError):
return "未知时间"
if stamp > 10_000_000_000:
stamp /= 1000
try:
return datetime.fromtimestamp(stamp).strftime("%Y-%m-%d %H:%M:%S")
except (OSError, OverflowError, ValueError):
return "未知时间"
def task_result_text(value):
if not value:
return "尚无结果"
try:
result = json.loads(value) if isinstance(value, str) else value
except json.JSONDecodeError:
return str(value)
if not isinstance(result, dict):
return str(result)
parts = []
for key, label in {
"code": "结果代码",
"status_code": "平台状态码",
"error": "错误",
"status": "状态",
"success": "是否成功",
"client_id": "客户端消息ID",
"server_id": "服务端消息ID",
}.items():
if key in result and result[key] not in (None, ""):
parts.append(f"{label}:{result[key]}")
message = result.get("message")
if isinstance(message, dict):
if message.get("content"):
parts.append("消息内容:" + str(message["content"]))
for key, label in (
("client_id", "客户端消息ID"),
("server_id", "服务端消息ID"),
):
if message.get(key):
parts.append(f"{label}:{message[key]}")
elif message:
parts.append("平台说明:" + str(message))
return ";".join(dict.fromkeys(parts)) or "平台业务结果已保存,可在底层数据库中定位"
class Backend(QThread):
snapshot_ready = Signal(dict)
response = Signal(str, bool, str)
@@ -211,8 +261,7 @@ class Window(QMainWindow):
("删除登录数据", self.erase),
("删除小号", self.delete_worker),
("历史事件操作", self.history_events),
("获取近期作品", lambda: self.fetch_works(False)),
("获取全部历史作品", lambda: self.fetch_works(True)),
("作品监控", self.fetch_works),
]:
self.button(actions, text, handler)
self.tabs.addTab(accounts, "账号组")
@@ -233,7 +282,7 @@ class Window(QMainWindow):
"目标 UID",
"动作",
"状态",
"结果(不含正文)",
"结果说明",
]
)
self.table.setSelectionBehavior(QTableWidget.SelectionBehavior.SelectRows)
@@ -343,17 +392,17 @@ class Window(QMainWindow):
return
self.send("history_fetch", {"id": account["id"]})
def fetch_works(self, all_pages):
def fetch_works(self):
account = self.selected()
if not account or account["role"] != "main" or not account["uid"]:
QMessageBox.information(self, "提示", "请先选择已登录并绑定身份的大号。")
return
self.send("works_all" if all_pages else "works_recent", {"id": account["id"]})
self.send("works_open", {"id": account["id"]})
def command_data(self, name, payload):
if name == "history_fetch":
self.history_dialog(payload)
elif name in ("works_recent", "works_all"):
elif name in ("works_open", "works_refresh"):
self.works_dialog(payload)
elif name == "history_enqueue":
QMessageBox.information(
@@ -373,7 +422,7 @@ class Window(QMainWindow):
def set_table_checks(table, state):
for row in range(table.rowCount()):
item = table.item(row, 0)
if item is not None:
if item is not None and not table.isRowHidden(row):
item.setCheckState(state)
@staticmethod
@@ -381,34 +430,56 @@ class Window(QMainWindow):
values = []
for row in range(table.rowCount()):
item = table.item(row, 0)
if item is not None and item.checkState() == Qt.CheckState.Checked:
if (
item is not None
and not table.isRowHidden(row)
and item.checkState() == Qt.CheckState.Checked
):
values.append(item.data(Qt.ItemDataRole.UserRole))
return values
def history_dialog(self, payload):
items = payload["items"]
labels = {
"digg": "点赞/作品互动",
"comment": "评论",
"follow": "关注",
"general_notice": "其他作品互动",
}
win = QDialog(self)
win.setWindowTitle(f"历史事件预览 — 共 {len(items)} 条(尚未创建任务)")
win.setWindowTitle(f"历史事件预览 — 当前可选 {len(items)} 条(尚未创建任务)")
win.resize(1500, 800)
layout = QVBoxLayout(win)
note = QLabel(
"勾选并确认后才按当前规则分配。历史任务优先级低于新消息;同组 UID 冷却、作品筛选和自操作保护继续生效。业务字段完整显示,认证凭据仍不输出。"
"列表来自本地缓存,并已按当前作品监控范围过滤。勾选并二次确认后才会产生真实关注或私信;新消息始终优先。完整平台数据只保存在本地数据库中。"
)
note.setWordWrap(True)
layout.addWidget(note)
table = QTableWidget(len(items), 10)
controls = QHBoxLayout()
controls.addWidget(QLabel("只看事件类型"))
kind_filter = QComboBox()
kind_filter.addItem("全部类型", "")
for kind in dict.fromkeys(item["kind"] for item in items):
kind_filter.addItem(labels.get(kind, "其他互动"), kind)
controls.addWidget(kind_filter)
select_all = QPushButton("全选当前显示")
clear = QPushButton("清空当前显示")
controls.addWidget(select_all)
controls.addWidget(clear)
controls.addStretch()
layout.addLayout(controls)
table = QTableWidget(len(items), 9)
table.setHorizontalHeaderLabels(
[
"选择",
"时间",
"类型",
"来源UID",
"来源昵称",
"发生时间",
"事件类型",
"来源用户",
"作品名称",
"作品ID",
"作品描述",
"评论/内容",
"评论或内容",
"通知ID",
"完整业务数据",
"数据来源",
]
)
for row, item in enumerate(items):
@@ -416,22 +487,25 @@ class Window(QMainWindow):
check.setFlags(Qt.ItemFlag.ItemIsEnabled | Qt.ItemFlag.ItemIsUserCheckable)
check.setCheckState(Qt.CheckState.Unchecked)
check.setData(Qt.ItemDataRole.UserRole, item["nid"])
check.setData(Qt.ItemDataRole.UserRole + 1, item["kind"])
table.setItem(row, 0, check)
users = []
for index, uid in enumerate(item["actor_uids"]):
name = (
item["actor_names"][index]
if index < len(item["actor_names"])
else ""
)
users.append(f"{name or '未知用户'}(UID {uid})")
values = [
item["create_time"],
{
"digg": "点赞/作品互动",
"comment": "评论",
"follow": "关注",
"general_notice": "其他作品互动",
}.get(item["kind"], item["kind"]),
", ".join(item["actor_uids"]),
", ".join(item["actor_names"]),
item["work_id"],
item["work_desc"],
item["comment"],
date_text(item["create_time"]),
labels.get(item["kind"], "其他互动"),
"、".join(users) or "未知用户",
item["work_desc"] or "未提供作品名称",
item["work_id"] or "无作品ID",
item["comment"] or "无文字内容",
item["nid"],
json.dumps(item["business"], ensure_ascii=False, default=str),
"实时收到" if item.get("origin") == "live" else "历史同步",
]
for column, value in enumerate(values, 1):
cell = QTableWidgetItem(str(value))
@@ -440,21 +514,29 @@ class Window(QMainWindow):
table.horizontalHeader().setSectionResizeMode(
QHeaderView.ResizeMode.ResizeToContents
)
table.horizontalHeader().setSectionResizeMode(9, QHeaderView.ResizeMode.Stretch)
table.horizontalHeader().setSectionResizeMode(4, QHeaderView.ResizeMode.Stretch)
layout.addWidget(table)
controls = QHBoxLayout()
select_all = QPushButton("全选")
clear = QPushButton("清空选择")
def apply_filter():
selected_kind = kind_filter.currentData()
for row in range(table.rowCount()):
item = table.item(row, 0)
table.setRowHidden(
row,
bool(selected_kind)
and (
item is None
or item.data(Qt.ItemDataRole.UserRole + 1) != selected_kind
),
)
kind_filter.currentIndexChanged.connect(apply_filter)
select_all.clicked.connect(
lambda: self.set_table_checks(table, Qt.CheckState.Checked)
)
clear.clicked.connect(
lambda: self.set_table_checks(table, Qt.CheckState.Unchecked)
)
controls.addWidget(select_all)
controls.addWidget(clear)
controls.addStretch()
layout.addLayout(controls)
buttons = QDialogButtonBox(
QDialogButtonBox.StandardButton.Ok | QDialogButtonBox.StandardButton.Cancel
)
@@ -465,7 +547,7 @@ class Window(QMainWindow):
return
selected = self.checked_table_data(table)
if not selected:
QMessageBox.information(self, "没有选择", "未创建任何历史任务。")
QMessageBox.information(self, "没有选择", "当前显示范围内未选择任何事件。")
return
if (
QMessageBox.question(
@@ -481,39 +563,47 @@ class Window(QMainWindow):
def works_dialog(self, payload):
items = payload["items"]
loaded = {item["aweme_id"] for item in items}
existing = set(payload["selected"])
hidden = existing - loaded
win = QDialog(self)
win.setWindowTitle(
("全部历史作品" if payload["all_pages"] else "近期作品")
+ f" — {len(items)} 条"
)
win.resize(1500, 800)
win.setWindowTitle(f"作品监控 — 本地已缓存 {len(items)} 个作品")
win.resize(1600, 800)
layout = QVBoxLayout(win)
all_mode = QCheckBox("监控全部作品(保持兼容;新作品自动包含)")
all_mode = QCheckBox("监控全部作品(以后发布的新作品也自动包含)")
all_mode.setChecked(payload["mode"] == "all")
layout.addWidget(all_mode)
note = QLabel(
"取消上方选项后,仅监控勾选作品的点赞、评论、收藏或其他带作品 ID 的互动;关注通知不受作品筛选。指定模式下新发布作品不会自动加入。"
+ (
f"\n当前还有 {len(hidden)} 个已选历史作品未包含在本次近期列表中,将继续保留;使用“获取全部历史作品”可统一管理。"
if hidden
else ""
)
"取消上方选项后,只监控表格中勾选作品的点赞、评论、收藏等互动;关注通知不受作品筛选。完整平台数据只保存在本地数据库中。"
)
note.setWordWrap(True)
layout.addWidget(note)
table = QTableWidget(len(items), 7)
schedule = QHBoxLayout()
schedule.addWidget(QLabel("运行期间自动刷新间隔"))
refresh_minutes = QSpinBox()
refresh_minutes.setRange(0, 10080)
refresh_minutes.setSpecialValueText("关闭")
refresh_minutes.setSuffix(" 分钟")
refresh_minutes.setValue(payload.get("refresh_interval", 3600) // 60)
schedule.addWidget(refresh_minutes)
refresh_now = QPushButton("强制增量刷新")
schedule.addWidget(refresh_now)
schedule.addWidget(
QLabel("首次打开会自动获取全部历史作品;停止账号组后不会定时请求。")
)
schedule.addStretch()
layout.addLayout(schedule)
table = QTableWidget(len(items), 10)
table.setHorizontalHeaderLabels(
[
"选择",
"发布时间",
"作品描述",
"作品ID",
"统计",
"封面URL",
"完整业务数据",
"点赞",
"评论",
"收藏",
"分享",
"封面地址",
"本地更新时间",
]
)
for row, item in enumerate(items):
@@ -526,13 +616,17 @@ class Window(QMainWindow):
)
check.setData(Qt.ItemDataRole.UserRole, item["aweme_id"])
table.setItem(row, 0, check)
statistics = item["statistics"]
values = [
item["create_time"],
item["desc"],
date_text(item["create_time"]),
item["desc"] or "未填写作品描述",
item["aweme_id"],
json.dumps(item["statistics"], ensure_ascii=False, default=str),
item["cover"],
json.dumps(item["business"], ensure_ascii=False, default=str),
statistics.get("digg_count", 0),
statistics.get("comment_count", 0),
statistics.get("collect_count", 0),
statistics.get("share_count", 0),
item["cover"] or "无封面地址",
date_text(item.get("updated")),
]
for column, value in enumerate(values, 1):
cell = QTableWidgetItem(str(value))
@@ -541,13 +635,13 @@ class Window(QMainWindow):
table.horizontalHeader().setSectionResizeMode(
QHeaderView.ResizeMode.ResizeToContents
)
table.horizontalHeader().setSectionResizeMode(6, QHeaderView.ResizeMode.Stretch)
table.horizontalHeader().setSectionResizeMode(2, QHeaderView.ResizeMode.Stretch)
table.setDisabled(all_mode.isChecked())
all_mode.toggled.connect(table.setDisabled)
layout.addWidget(table)
controls = QHBoxLayout()
select_all = QPushButton("全选当前列表")
clear = QPushButton("清空当前列表")
select_all = QPushButton("全选")
clear = QPushButton("清空选择")
select_all.clicked.connect(
lambda: self.set_table_checks(table, Qt.CheckState.Checked)
)
@@ -562,15 +656,21 @@ class Window(QMainWindow):
QDialogButtonBox.StandardButton.Save
| QDialogButtonBox.StandardButton.Cancel
)
buttons.button(QDialogButtonBox.StandardButton.Save).setText("保存设置")
buttons.button(QDialogButtonBox.StandardButton.Cancel).setText("取消")
buttons.accepted.connect(win.accept)
buttons.rejected.connect(win.reject)
layout.addWidget(buttons)
def force_refresh():
win.reject()
self.send("works_refresh", {"id": payload["account"]})
refresh_now.clicked.connect(force_refresh)
if win.exec() != QDialog.DialogCode.Accepted:
return
mode = "all" if all_mode.isChecked() else "selected"
selected = list(hidden)
selected.extend(self.checked_table_data(table))
selected = list(dict.fromkeys(selected))
selected = self.checked_table_data(table)
if mode == "selected" and not selected:
QMessageBox.warning(
self,
@@ -579,7 +679,13 @@ class Window(QMainWindow):
)
return
self.send(
"work_filter", {"id": payload["account"], "mode": mode, "ids": selected}
"work_filter",
{
"id": payload["account"],
"mode": mode,
"ids": selected,
"refresh_interval": refresh_minutes.value() * 60,
},
)
def selected(self):
@@ -724,6 +830,7 @@ class Window(QMainWindow):
"cooldown": cooldown.value() * 60,
"work_mode": rule.get("work_mode", "all"),
"work_ids": rule.get("work_ids", []),
"works_refresh_interval": rule.get("works_refresh_interval", 3600),
}
try:
validate_rule(value)
@@ -995,7 +1102,7 @@ class Window(QMainWindow):
task["target"],
"关注" if task["action"] == "follow" else "私信",
labels[task["status"]],
task["result"],
task_result_text(task["result"]),
]
for col, value in enumerate(values):
self.table.setItem(row, col, QTableWidgetItem(str(value)))
+84 -9
View File
@@ -218,6 +218,66 @@ class StoreTest(unittest.TestCase):
{states[nid] for nid in ("1011", "1013", "1014")}, {"dispatched"}
)
def test_persistent_caches_merge_incrementally_filter_and_hide_credentials(self):
self.store.cache_works(
self.main,
[
{
"aweme_id": "700",
"create_time": 100,
"desc": "旧描述",
"statistics": {"digg_count": 1},
"cover": "https://example.invalid/cover",
"business": {"Authorization": "secret", "desc": "旧描述"},
}
],
)
self.store.cache_works(
self.main,
[
{
"aweme_id": "700",
"create_time": 100,
"desc": "新描述",
"statistics": {"digg_count": 2},
"cover": "https://example.invalid/cover",
"business": {"Authorization": "secret", "desc": "新描述"},
}
],
)
selected = work_notice("101", "3001", "901", "700")
other = work_notice("101", "3002", "902", "701")
followed = notice("101", "3003", "903")
selected["Authorization"] = "secret"
self.store.cache_notices(self.main, [selected, other, followed], "history")
self.store.cache_notices(self.main, [selected], "live")
self.store.set_work_filter(self.main, "selected", ["700"], 0)
self.store.close()
self.store = Store(self.temp.name)
works = self.store.cached_works(self.main)
self.assertEqual(len(works), 1)
self.assertEqual(works[0]["desc"], "新描述")
self.assertEqual(works[0]["statistics"]["digg_count"], 2)
work_business = json.loads(
self.store.db.execute(
"SELECT business FROM cached_works WHERE source=? AND aweme_id=?",
(self.main, "700"),
).fetchone()["business"]
)
self.assertEqual(work_business["Authorization"], "[凭据已隐藏]")
cached = {item["nid"]: item for item in self.store.cached_notices(self.main)}
self.assertEqual(set(cached), {"3001", "3003"})
self.assertEqual(cached["3001"]["origin"], "live")
business = self.store.cached_notice(self.main, "3001")
assert business is not None
self.assertEqual(business["Authorization"], "[凭据已隐藏]")
self.assertEqual(
json.loads(self.store.account(self.main)["rule"])["works_refresh_interval"],
0,
)
def test_live_tasks_preempt_history_and_promote_same_notification(self):
self.store.set_rule(self.main, rule(dm=False, cooldown=0))
self.store.ingest(self.main, notice("101", "2001", "901"), origin="history")
@@ -657,12 +717,18 @@ class EngineTest(unittest.IsolatedAsyncioTestCase):
await self.engine.start([ident])
self.assertEqual(self.engine.desired, set())
async def test_im_result_classification_keeps_business_data_and_hides_credentials(self):
async def test_im_result_classification_keeps_business_data_and_hides_credentials(
self,
):
session = AsyncMock()
task = {"params": json.dumps(rule()), "action": "dm", "target": "999"}
session.im.return_value = {
"success": True,
"message": {"client_id": "c1", "server_id": "s1", "content": "完整私信正文"},
"message": {
"client_id": "c1",
"server_id": "s1",
"content": "完整私信正文",
},
"token": "AUTH-CREDENTIAL",
}
status, result = await self.engine.perform(session, task)
@@ -939,18 +1005,27 @@ class EngineTest(unittest.IsolatedAsyncioTestCase):
"history_enqueue", {"id": self.main, "ids": ["3001"]}
)
self.assertEqual(result, {"selected": 1, "events": 1, "tasks": 2})
fetched = await self.engine.command("works_recent", {"id": self.main})
self.assertEqual(fetched["items"], works)
self.assertFalse(fetched["all_pages"])
fetched = await self.engine.command("works_open", {"id": self.main})
self.assertEqual(fetched["items"][0]["aweme_id"], "700")
self.assertEqual(fetched["refresh_interval"], 3600)
saved = await self.engine.command(
"work_filter", {"id": self.main, "mode": "selected", "ids": ["700"]}
"work_filter",
{
"id": self.main,
"mode": "selected",
"ids": ["700"],
"refresh_interval": 120,
},
)
self.assertEqual(
saved, {"mode": "selected", "count": 1, "refresh_interval": 120}
)
self.assertEqual(saved, {"mode": "selected", "count": 1})
rule_value = json.loads(self.engine.store.account(self.main)["rule"])
self.assertEqual(rule_value["work_mode"], "selected")
self.assertEqual(rule_value["work_ids"], ["700"])
session.works.assert_awaited_once_with(False)
session.history_notices.assert_awaited_once()
self.assertEqual(rule_value["works_refresh_interval"], 120)
session.works.assert_awaited_once_with(True, None)
session.history_notices.assert_awaited_once_with(set())
async def test_move_during_identity_check_does_not_cross_groups(self):
other = self.engine.store.add("main", "另一组", "102")
+19 -9
View File
@@ -150,14 +150,12 @@ class OnboardingUiTest(unittest.TestCase):
patch.object(self.window, "send") as send,
):
self.window.history_events()
self.window.fetch_works(False)
self.window.fetch_works(True)
self.window.fetch_works()
self.assertEqual(
[call.args for call in send.call_args_list],
[
("history_fetch", {"id": "main"}),
("works_recent", {"id": "main"}),
("works_all", {"id": "main"}),
("works_open", {"id": "main"}),
],
)
@@ -174,6 +172,7 @@ class OnboardingUiTest(unittest.TestCase):
"actor_uids": ["999"],
"actor_names": ["完整昵称"],
"comment": "完整评论正文",
"origin": "history",
"business": {
"nid_str": "9007199254740993",
"comment": "完整评论正文",
@@ -184,8 +183,9 @@ class OnboardingUiTest(unittest.TestCase):
def choose(dialog):
table = dialog.findChild(QTableWidget)
assert (
table is not None and table.item(0, 9).text().find("完整评论正文") >= 0
assert table is not None and "完整评论正文" in table.item(0, 6).text()
assert "{" not in " ".join(
table.item(0, column).text() for column in range(table.columnCount())
)
table.item(0, 0).setCheckState(Qt.CheckState.Checked)
return QDialog.DialogCode.Accepted
@@ -207,7 +207,7 @@ class OnboardingUiTest(unittest.TestCase):
"account": "main",
"mode": "selected",
"selected": [],
"all_pages": True,
"refresh_interval": 3600,
"items": [
{
"aweme_id": "700",
@@ -215,6 +215,7 @@ class OnboardingUiTest(unittest.TestCase):
"create_time": 123,
"statistics": {"collect_count": 8},
"cover": "https://example.invalid/full.jpg",
"updated": 456,
"business": {"aweme_id": "700", "desc": "完整作品描述"},
}
],
@@ -222,7 +223,10 @@ class OnboardingUiTest(unittest.TestCase):
def choose(dialog):
table = dialog.findChild(QTableWidget)
assert table is not None and "完整作品描述" in table.item(0, 6).text()
assert table is not None and "完整作品描述" in table.item(0, 2).text()
assert "{" not in " ".join(
table.item(0, column).text() for column in range(table.columnCount())
)
table.item(0, 0).setCheckState(Qt.CheckState.Checked)
all_mode = next(
box
@@ -238,7 +242,13 @@ class OnboardingUiTest(unittest.TestCase):
):
self.window.works_dialog(payload)
send.assert_called_once_with(
"work_filter", {"id": "main", "mode": "selected", "ids": ["700"]}
"work_filter",
{
"id": "main",
"mode": "selected",
"ids": ["700"],
"refresh_interval": 3600,
},
)
def test_anonymous_profile_is_never_bound(self):