diff --git a/WINDOWS-README.txt b/WINDOWS-README.txt index 779bc45..c225ada 100644 --- a/WINDOWS-README.txt +++ b/WINDOWS-README.txt @@ -1,7 +1,7 @@ -抖音账号助手 0.1.6 — Windows 10/11 x64 +抖音账号助手 0.1.7 — Windows 10/11 x64 【安装】 -只需运行 DouyinAccounts-0.1.6-Windows-x64-Setup.exe。 +只需运行 DouyinAccounts-0.1.7-Windows-x64-Setup.exe。 无需安装 Python、Qt、Patchright、Node.js 或另行下载浏览器。 安装程序不会自动启动业务;没有开机自启、自动登录或默认私信。 默认使用随包的 fingerprint-chromium(adryfish 项目)和 Patchright 驱动,不再回退到系统 Chrome。 @@ -17,6 +17,7 @@ 若关闭了浏览器,可点击“登录 / 刷新信息”继续;已绑定账号点击此按钮可刷新资料和统计。 只保留一个抖音首页或自己的页面;多个匹配标签会拒绝执行,防止选错。 3. 选择大号,配置规则。默认关闭;可配置通知类别、关注、私信、正文、不同任务批次间隔和同组 UID 冷却。 + 作品监控默认“全部作品”。点击“获取近期作品”或“获取全部历史作品”,可切换为只监控勾选作品。 冷却默认 240 分钟(4 小时),设为 0 可关闭;勾选大号后点击启动。规则未开启不会产生或执行业务任务。 4. 每条匹配通知只给一个关联小号,按组内轮询,不跨组操作。 一条聚合通知可能包含多个来源用户,每个有效来源用户各执行配置动作。 @@ -37,6 +38,26 @@ 此操作不删除账号定义、任务和历史,也不影响其他账号。 更换 Chrome 可执行路径会使用不同登录目录,不复制 Cookie。 +【历史事件批量操作(0.1.7 新增)】 +选中已绑定身份的大号,点击“历史事件操作”。程序通过网页原生只读接口分页获取平台当前可返回的全部互动通知,固定 is_mark_read=0,不会自动标记已读。 +获取完成后显示时间、类型、来源 UID/昵称、作品 ID/描述、评论内容、通知 ID和完整业务数据;此时尚未创建任务。 +可逐条勾选、全选或取消。点击确认后还会出现第二次真实操作警告;只有再次确认才按当前规则、作品范围、冷却和自操作保护创建任务。 +历史任务优先级低于新消息:Worker 领取任务时先取新消息;同 UID 新消息会取消尚未执行的历史任务。已经执行、正在执行或待核对的历史任务不会被静默改写或重发。 +同一通知如果先由历史列表入队、随后实时推送到达,尚未执行的任务会提升为新消息优先级,不重复创建。 +重复历史通知、规则未开启、通知类别未勾选、作品范围外、无来源用户、冷却中或目标是执行号自身时都可能显示“选择成功但任务数较少”,具体原因查看运行记录。 + +【指定作品监控(0.1.7 新增)】 +选中已绑定身份的大号,可点击“获取近期作品”(一页)或“获取全部历史作品”(分页至平台返回结束)。列表显示发布时间、完整描述、作品 ID、统计、封面 URL和完整业务数据。 +默认“监控全部作品”,保持升级前行为并自动包含以后发布的新作品。取消该选项后进入“指定作品”模式,至少勾选一个作品;此后新作品不会自动加入。 +指定作品模式按通知内的作品 ID统一过滤点赞、评论、收藏或其他作品互动。平台当前没有暴露稳定的独立收藏通知常量,因此不会硬编码未经证实的类型;只要通知携带作品 ID,就使用同一筛选规则。 +关注通知没有作品 ID,不受作品筛选。近期列表不会丢弃未显示的既有历史选择;需要统一删除旧选择时使用“获取全部历史作品”。 +作品选择持久化在大号规则中,新通知和用户确认的历史事件使用同一范围。已入队任务继续使用原规则快照,不因修改作品范围而静默改写。 + +【业务输出与认证凭据(0.1.7 调整)】 +UID、通知 ID、昵称、作品 ID/描述、评论内容、私信正文、URL、平台资料和结构化业务结果均按原文显示和写入本地日志,不再使用哈希、尾号或业务字段隐藏。 +Cookie、Authorization、Session、Token、密码、verifyFp、a_bogus、X-Bogus及其他签名/认证字段仍强制隐藏,避免按日日志成为可直接接管账号的凭据文件。 +完整业务输出可能包含个人信息并显著增加日志体积。日志仅保存在本机;请限制数据目录访问并自行确定保留期限。 + 【同组 UID 冷却与双动作并行(0.1.6 新增)】 冷却按大号组共享:该组任一小号开始操作某 UID 后,组内其他小号在设定时长内也不会再次操作该 UID;不同大号组互不影响。 冷却从任务实际领取、即将请求平台时开始并持久化,程序重启不重置。只排队后因停止、切换归属或删除小号而取消的任务不会启动冷却。 @@ -52,7 +73,7 @@ 记录推送接收与去重、详情请求/重试、规则不匹配、无执行小号、自操作保护、UID 冷却、并行批次、轮询分配、任务入队/执行/结果、旧任务依赖等待、人工核对、账号操作和异常。 例如“触发者就是执行小号自身”会明确显示跳过原因;“分配完成”不代表已经关注或发送私信。 监听无新通知时每分钟写一次心跳,重复等待状态不反复刷屏。任务列表也仅展示最近 2000 条,完整历史仍在数据库中。 -日志不保存 Cookie、签名、私信正文、原始请求/响应;通知标识使用稳定摘要,目标 UID 只显示尾号,账号用本地名称标识。 +日志保留完整业务字段,包括私信正文、通知 ID、目标 UID、作品资料、URL和业务结果;Cookie、Token、密码、Authorization和签名字段仍不保存。 磁盘或权限异常会在 UI 明确提示;该期间仅有最近 2000 条内存记录,缺失的文件日志不会自动补写,请及时检查。 旧版 activity.log 及备份仍保留,未记录过的旧过程不会补造。日志会持续占用磁盘,请自行备份管理。 @@ -70,7 +91,7 @@ 卸载/覆盖升级默认保留账号数据。从 0.1.0 升级旧数据库时会先备份再迁移,保留原 UID、归属、任务和历史。 升级 0.1.2 前请先停止业务并关闭旧 Chrome。新内核使用独立登录目录,需要自行重新登录原账号。 旧目录及旧浏览器配置保留,但新版本不自动沿用旧 Chrome 的路径配置;如需自定义请重新设置指纹浏览器路径。 -0.1.3 及更新版本首次启动旧库时会增加通知重试字段;0.1.5 增加删除标记字段;0.1.6 增加组内 UID 冷却表和规则字段。旧账号、通知 ID、任务和登录目录不变。 +0.1.3 及更新版本首次启动旧库时会增加通知重试字段;0.1.5 增加删除标记字段;0.1.6 增加组内 UID 冷却表;0.1.7 增加事件来源和任务优先级字段,并在大号规则增加作品范围。旧账号、通知 ID、任务和登录目录不变。 0.1.6 将旧规则的“必须关注成功后再私信”自动关闭;旧版已排队且带前置依赖的任务仍按原快照安全处理,不重写历史。 0.1.4 修复了新点赞、评论推送后无法取得详情的问题:复用网页自身请求层,自动补齐当前运行时参数,保持原始 ID 精度。 本次已经实际验证“新增点赞和评论 → 两条推送 → 两条完整详情”,未执行自动关注或私信。 diff --git a/build_windows.py b/build_windows.py index 87d2a5b..1bd4eed 100644 --- a/build_windows.py +++ b/build_windows.py @@ -168,7 +168,7 @@ def main(): break source_hashes = {p.as_posix(): sha256(p) for p in sorted(Path("src").glob("*.py"))} manifest = { - "app_version": "0.1.6", + "app_version": "0.1.7", "python": sys.version, "chrome": browser, "dependencies": versions, @@ -208,7 +208,7 @@ def main(): "onedir 和离线冒烟已完成。安装 Inno Setup 6 后传 --iscc PATH 以生成最终单文件安装包。" ) subprocess.run([str(iscc), str(ROOT / "installer.iss")], check=True) - installer = ROOT / "release/DouyinAccounts-0.1.6-Windows-x64-Setup.exe" + installer = ROOT / "release/DouyinAccounts-0.1.7-Windows-x64-Setup.exe" (installer.parent / "SHA256SUMS.txt").write_text( sha256(installer) + " " + installer.name + "\n", encoding="ascii" ) diff --git a/docs/2026-09-07-11-32-history-events-and-work-monitoring.md b/docs/2026-09-07-11-32-history-events-and-work-monitoring.md new file mode 100644 index 0000000..05eb0cc --- /dev/null +++ b/docs/2026-09-07-11-32-history-events-and-work-monitoring.md @@ -0,0 +1,327 @@ +# 0.1.7:历史事件预览操作、指定作品监控与完整业务输出 + +## 需求与确认结果 + +本次需求包括: + +1. 系统输出不再对业务数据脱敏。 +2. 大号可以获取全部历史事件,预览、选择后交给小号执行;新消息优先。 +3. 大号可以获取近期或全部历史作品,选择仅监控指定作品的点赞、收藏、评论等互动。 + +通过交互确认: + +- UID、通知 ID、昵称、评论、私信正文、作品 ID/描述、URL、平台资料和结构化业务结果按原文输出。 +- Cookie、Authorization、Session、Token、密码和签名继续强制隐藏,避免日志成为可直接接管账号的凭据文件。 +- 历史事件必须先完整获取和预览,用户逐条勾选或全选,二次确认后才能产生真实任务。 +- 作品监控升级后默认“全部作品”,保持旧行为;用户主动保存选择后进入“指定作品”模式,新发布作品不会自动加入。 + +排查和验证期间先通过 Windows MCP“停止全部”,保留用户浏览器登录。所有真实账号接口调研均为 GET/SDK 只读请求,固定 `is_mark_read=0`,没有关注、私信、标记已读或自动登录。 + +## 只读接口调研 + +### 历史互动通知 + +在已登录且身份核验一致的页面中,通过 browser-harness 检查网页现有 Webpack SDK: + +- 通知模块当前包含通知列表、计数、详情、点赞用户列表和相关只读函数。 +- “全部互动”使用 `notice_group=700`。 +- 列表请求使用 `count`、`min_time`、`max_time` 和 `has_more` 分页。 +- 本次现场以每页 50 条连续读取前 5 页,得到 250 个唯一通知 ID,接口仍报告有后续页;因此产品不能只取固定 5 页。 +- 现场历史类型包含评论、关注和点赞/作品互动;行内包含稳定通知 ID、时间、作品 ID及嵌套用户、评论、点赞数据。 +- 所属 `user_id` 可与已绑定大号逐条核验。 + +产品实现持续分页直到平台返回 `has_more=0`,并检测重复/不推进游标。设置 1000 页、50000 条的安全上限;超过时明确报错,不把截断结果冒充“全部”。 + +通知列表的网页函数默认可能设置已读参数,因此产品显式传入 `is_mark_read=0`,并在离线浏览器 fixture 中强制断言该参数。 + +### 作品列表 + +网页现有作品函数使用 `/aweme/v1/web/aweme/post/`,请求字段包括: + +- `sec_user_id`; +- `max_cursor`; +- `count`; +- 时间列表选项。 + +响应包含: + +- `aweme_id`; +- 完整描述; +- 发布时间; +- 作者资料; +- 点赞、评论、收藏等 statistics; +- 视频封面或图集 URL; +- `has_more` 和 `max_cursor`。 + +“近期作品”取一页 18 条;“全部历史作品”持续翻页直到结束,检测重复游标,安全上限为 1000 页/18000 条,超限明确失败。 + +### 收藏事件边界 + +当前页面代码和现场 250 条通知中没有发现稳定、独立的“收藏通知类型常量”。点赞通知存在多个子类型,其中现场出现的子类型 22 实际可对应图文作品点赞,不能错误硬编码为收藏。 + +因此产品采用可验证的结构策略: + +- 若通知含 `favorite` 或 `collect` 业务对象,提取其中的来源用户和作品。 +- 无论事件被平台归类为点赞、评论、收藏或其他,只要携带作品 ID,就执行同一作品范围过滤。 +- UI 将 `general_notice` 显示为“收藏/其他作品互动”。 +- 不声称平台一定会为每次收藏推送通知;系统只能处理平台实际提供的推送或历史记录。 + +## 网页原生只读请求 + +新增通用 `native_read_script()`,继续复用网页自身 Axios/XHR 请求层,不手工拼接简化 fetch: + +1. 动态取得当前页面 Webpack require。 +2. 以 JSON 字符串形式的**完整端点字面量**匹配模块和导出函数,避免 `/notice/count/`、`/notice/delete/` 等前缀相似端点误匹配。 +3. 历史通知函数还要求源码含 `is_new_notice` 标记。 +4. 作品函数使用网页现有 `fetchUserPost` 导出。 +5. 临时 Axios response observer 只匹配同源、精确 pathname及本次分页参数。 +6. 从 XHR `responseText` 读取原始文本,再交给 Python `json.loads`,避免 JavaScript 对 64 位 ID舍入。 +7. 成功、失败和超时都在 `finally` 中删除 observer。 +8. 不修改全局 XHR、不篡改响应、不创建额外登录状态。 + +最初通用匹配只检查端点前缀,现场只读验证误选到 notice count 函数并返回业务码 5。修复为完整端点字面量后,再次现场验证: + +```json +{ + "identity_verified": true, + "history_http": 200, + "history_status": 0, + "history_count": 3, + "history_mark_read": 0, + "history_ids_exact": true, + "history_identity_verified": true, + "works_http": 200, + "works_status": 0, + "works_count": 3, + "works_ids_exact": true, + "works_identity_verified": true, + "real_write_actions": 0 +} +``` + +调研 SSH/CDP 临时转发随后关闭,不关闭用户浏览器。 + +## 历史事件操作 + +### UI 流程 + +大号工具栏新增“历史事件操作”: + +1. 只允许已登录且已绑定身份的大号。 +2. 点击后只读获取全部可用历史互动。 +3. 完成后打开预览表,显示: + - 时间; + - 类型; + - 来源 UID和昵称; + - 作品 ID和完整描述; + - 评论/内容; + - 通知 ID; + - 完整业务 JSON。 +4. 初始不勾选,不产生任务。 +5. 支持逐条勾选、全选和清空。 +6. 点击确定后再次显示“可能产生真实关注和私信”的确认框,默认按钮为“否”。 +7. 用户再次确认后,UI只发送已加载通知 ID;Engine根据内存缓存取回原始对象,拒绝未加载或伪造 ID。 + +缓存仅存在于当前程序运行期;重启后必须重新获取,避免使用陈旧平台结果。 + +### 入库和去重 + +`Store.ingest(source, notice, origin)` 新增来源: + +- `live`:实时推送; +- `history`:用户确认的历史列表。 + +事件仍以 `(source,nid)` 唯一,历史列表和实时推送不会对同一通知创建两套事件。历史确认使用当前大号规则快照,并继续经过: + +- 规则开关; +- 通知类别; +- 指定作品范围; +- 有效来源 UID; +- 自操作保护; +- 组共享 UID 冷却; +- pending/running/unknown 重复保护。 + +界面确认结果显示:选择数量、新增事件数量和实际创建任务数量。三者可能不同,详细原因写入每日运行日志。 + +## 新消息优先 + +数据库迁移新增: + +```sql +events.origin TEXT NOT NULL DEFAULT 'live'; +tasks.priority INTEGER NOT NULL DEFAULT 100; +``` + +- 实时任务优先级 100。 +- 历史任务优先级 0。 +- Worker领取时 `ORDER BY priority DESC,id`。 +- waiting events也先分配 live,再分配 history。 +- 正在执行的历史任务不强行中断;完成当前平台请求后再领取实时任务。 + +### 同 UID冲突 + +若实时新消息到达时,同一大号组和目标 UID已有尚未执行的历史任务: + +- pending历史任务改为 cancelled; +- 记录“新消息优先,取消未执行历史任务”; +- 创建实时任务; +- 不把历史任务改派或重发。 + +若历史任务已 running 或 unknown,继续遵守“在途不强停、结果不确定不重发”,不创建可能重复打扰的新任务。 + +若同一通知先从历史列表进入、随后相同 nid实时推送到达: + +- events.origin 从 history 提升为 live; +- 其 pending任务优先级提升到 100; +- 不重复创建事件或任务。 + +## 指定作品监控 + +大号工具栏新增: + +- “获取近期作品”; +- “获取全部历史作品”。 + +作品表显示发布时间、完整描述、作品 ID、statistics、封面 URL和完整业务 JSON。 + +规则新增: + +```json +{ + "work_mode": "all", + "work_ids": [] +} +``` + +### 全部作品模式 + +- 升级默认值。 +- 与旧版本一致。 +- 以后发布的新作品自动包含。 + +### 指定作品模式 + +- 至少选择一个作品。 +- `work_ids` 去重并持久化,最多 50000 个 ASCII 十进制 ID。 +- 新作品不会自动加入。 +- 近期列表没有显示的既有历史选择会保留,避免用户只刷新近期作品就误删旧选择。 +- 需要统一移除旧选择时,使用“获取全部历史作品”。 +- 点赞、评论、收藏/其他作品互动必须提取出作品 ID且在列表内,否则事件记录为 ignored、不生成任务。 +- 关注通知没有作品 ID,不受作品范围限制。 +- 修改范围只影响之后入库的事件;已创建任务保留原规则快照。 + +`Window.rules()` 保存普通规则时也会携带原 `work_mode/work_ids`,不会因为修改私信正文或间隔而意外重置作品选择。 + +## 完整业务输出与凭据保护 + +### 已移除的业务字段隐藏 + +以下内容现在原样进入 UI、任务结果、本地 profile和每日日志: + +- 完整 UID和通知 ID; +- 昵称、平台资料和手机号等资料字段; +- 作品 ID、描述、统计、封面 URL; +- 评论内容; +- 私信正文; +- 业务 URL; +- SDK成功、失败和待核对结果结构。 + +删除了通知 ID哈希、UID尾号显示、URL隐藏、正文隐藏和“复杂数据不记录”。结构化 dict/list 使用 UTF-8 JSON输出,不截断单行业务文本。 + +`get_current_user` 和异步 Session现在返回完整业务资料,并补充兼容用的 `douyin_id/avatar_url` 字段。Store保存完整业务 profile。 + +### 仍强制保护的认证字段 + +以下字段无论嵌套深度都替换为 `[凭据已隐藏]`: + +- Cookie; +- Authorization; +- Session ID; +- Token/msToken; +- password/secret; +- verifyFp; +- a_bogus、X-Bogus、X-TT-Params和签名。 + +URL本身保留,但认证查询参数值会隐藏。控制字符、换行和双向文本控制符继续清理,防止伪造日志行;这不是业务字段脱敏。 + +任务结果以前只保留私信 client/server ID。本次改为保存完整 business result;若结果中包含认证键,先递归替换凭据值,再写 SQLite和日志。未知结果继续附加 `SDK_RESULT_UNCONFIRMED`,不会自动重发。 + +## 测试 + +### Python测试 + +Linux与Windows完整 pytest: + +```text +73 passed +``` + +新增覆盖: + +- 历史通知双页分页、完整 64 位 nid、`is_mark_read=0`; +- 作品双页分页、完整 aweme ID、描述、统计和封面 URL; +- 历史/作品重复游标拒绝; +- 历史或作品所属身份不符时失败关闭; +- Native SDK模块精确端点匹配、Axios observer成功/失败后清理; +- 历史预览尚未确认时任务数为 0; +- 未加载的历史 ID拒绝; +- 用户确认后按规则创建任务; +- 指定作品允许所选点赞/收藏结构,拒绝范围外作品; +- 关注不受作品筛选; +- live任务先于history任务领取; +- 同 UID实时消息取消未执行历史任务; +- 同 nid历史任务提升为实时优先级; +- 默认作品模式、作品 ID类型/数量验证和旧规则迁移; +- UI按钮只接受已绑定大号; +- 历史预览业务字段原样显示、默认不选、二次确认; +- 作品预览保存指定模式和完整业务字段; +- profile保留手机号等业务字段,Cookie值隐藏; +- 私信成功结果保留完整正文,Token值隐藏。 + +### Windows fingerprint-chromium矩阵 + +新增: + +```text +native_history_and_works_pagination_identity_precision_and_cleanup +``` + +真实离线 Chromium fixture要求: + +- 历史请求 `notice_group=700`、`is_mark_read=0`; +- 历史分页两页; +- 作品使用绑定账号的 sec_uid并分页两页; +- 通知和作品 64 位 ID保留字符串精度; +- 每条用户/作者身份匹配; +- 所有临时 Axios observer清理,不破坏原有 observer。 + +原有详情、登录、身份隔离、目录/端口、固定种子、重连、小号删除和浏览器关闭矩阵继续通过。报告 `real_write_actions=0`。 + +### 构建 + +- 13 个相关 Python文件 LSP error检查:0。 +- Windows pytest:73项通过。 +- 打包后 smoke:Qt、核心模块、SQLite、Patchright、内置浏览器路径均通过;accounts=0、real_write_actions=0。 +- 23个 Python源文件 SHA256与Windows build-manifest全部一致。 +- 浏览器和运行时保持 fingerprint-chromium 148.0.7778.215、Patchright 1.62.3、Python 3.12.10。 + +## 发布 + +```text +C:\Users\rogee\Desktop\抖音账号助手-0.1.7\ + DouyinAccounts-0.1.7-Windows-x64-Setup.exe + SHA256SUMS.txt + 使用说明.txt +``` + +- 安装包大小:222682833字节。 +- SHA256:`5401c81ab96565e41ea4ffdaf423723faafae4f4088f39ca6967afac750a8428`。 + +未自动安装、启动业务或把真实历史事件入队。建议覆盖安装后: + +1. 先打开“获取近期作品”,确认完整业务字段和默认“全部作品”。 +2. 如需节省小号配额,切换为指定作品并保存。 +3. 打开“历史事件操作”,先选少量可验证事件,确认预览和二次确认流程。 +4. 手动启动账号组,观察新消息是否优先于历史待执行任务。 + +历史批量确认会产生真实关注/私信,必须由用户自行验收和承担平台行为结果。 diff --git a/installer.iss b/installer.iss index 8bdf5a3..76a79ef 100644 --- a/installer.iss +++ b/installer.iss @@ -2,14 +2,14 @@ [Setup] AppId={{C9C3EE3B-3666-4E41-A52B-BB309838F701} AppName=抖音账号助手 -AppVersion=0.1.6 +AppVersion=0.1.7 DefaultDirName={localappdata}\Programs\DouyinAccounts DefaultGroupName=抖音账号助手 PrivilegesRequired=lowest ArchitecturesAllowed=x64compatible ArchitecturesInstallIn64BitMode=x64compatible OutputDir=release -OutputBaseFilename=DouyinAccounts-0.1.6-Windows-x64-Setup +OutputBaseFilename=DouyinAccounts-0.1.7-Windows-x64-Setup Compression=lzma2 SolidCompression=yes WizardStyle=modern diff --git a/src/account_engine.py b/src/account_engine.py index b7e3b46..2fcdd08 100644 --- a/src/account_engine.py +++ b/src/account_engine.py @@ -5,12 +5,37 @@ import logging import time from account_browser import BrowserManager, validate_config -from account_log import DailyLog +from account_log import DailyLog, visible from account_session import Session, SessionError -from account_store import Store, decode +from account_store import Store, decode, notice_kind, notice_work_id 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) @@ -31,6 +56,8 @@ 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( "历史摘要", @@ -279,16 +306,13 @@ class Engine: rule = decode(task["params"]) if task["action"] == "follow": result = await session.follow(task["target"]) - return result.get("status", "unknown"), result + return result.get("status", "unknown"), visible(result) result = await session.im( "send", task["target"], text=rule["text"], confirm=True ) + output = visible(result) if type(result.get("success")) is bool and result["success"]: - message = result.get("message") or {} - return "succeeded", { - "client_id": message.get("client_id"), - "server_id": message.get("server_id"), - } + return "succeeded", output known = { "WRONG_ORIGIN", "LOGIN_CHECK_FAILED", @@ -300,7 +324,7 @@ class Engine: "MESSAGE_BUILD_FAILED", } if result.get("error") in known: - return "failed", {"code": result["error"]} + return "failed", output # Explicit SDK rejection is a failed attempt; missing/ambiguous results are never resent. if ( type(result.get("success")) is bool @@ -308,8 +332,10 @@ class Engine: and type(result.get("status_code")) is int and result["status_code"] != 0 ): - return "failed", {"code": result["status_code"]} - return "unknown", {"code": "SDK_RESULT_UNCONFIRMED"} + return "failed", output + if isinstance(output, dict): + output.setdefault("code", "SDK_RESULT_UNCONFIRMED") + return "unknown", output async def perform_batch(self, session, tasks): results = await asyncio.gather( @@ -387,6 +413,110 @@ class Engine: session = None await self.pause(ident, 5) + async def fetch_history(self, ident): + account = self.store.account(ident) + if account["role"] != "main" or account["uid"] is None: + raise ValueError("请选择已绑定身份的大号") + session = await self.open(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), + ) + return {"account": ident, "items": previews} + + def enqueue_history(self, ident, ids): + if ( + not isinstance(ids, list) + or len(ids) > 50000 + or any( + not isinstance(item, str) or not item.isascii() or not item.isdecimal() + for item in ids + ) + ): + 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("选择包含未加载的历史事件") + 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 + ) + tasks = ( + self.store.db.execute("SELECT count(*) FROM tasks").fetchone()[0] - before + ) + self.store.log( + "历史确认", + "用户已确认选择;历史事件按低优先级分配,新消息始终优先", + source=ident, + 选择数=len(ids), + 新增事件数=added, + 新增任务数=tasks, + ) + return {"selected": len(ids), "events": added, "tasks": tasks} + + async def fetch_works(self, ident, all_pages): + 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 "近期", + ) + rule = decode(account["rule"]) + return { + "account": ident, + "items": works, + "mode": rule["work_mode"], + "selected": rule["work_ids"], + "all_pages": all_pages, + } + + def save_work_filter(self, ident, mode, ids): + 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} + async def command(self, name, data): labels = { "add": "添加账号", @@ -401,6 +531,11 @@ class Engine: "close_browser": "关闭浏览器", "erase": "清除登录数据", "delete_worker": "删除小号", + "history_fetch": "获取全部历史事件", + "history_enqueue": "确认历史事件操作", + "works_recent": "获取近期作品", + "works_all": "获取全部历史作品", + "work_filter": "保存作品监控范围", } label = labels.get(name, "未知操作") self.store.log("操作请求", label, account=data.get("id")) @@ -439,6 +574,14 @@ class Engine: elif name == "open": self.store.account(data["id"]) self.begin_login(data["id"]) + elif name == "history_fetch": + 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 == "work_filter": + return self.save_work_filter(data["id"], data.get("mode"), data.get("ids")) elif name == "move": self.store.move(data["id"], data.get("owner")) owner = data.get("owner") diff --git a/src/account_log.py b/src/account_log.py index 00604dd..a337635 100644 --- a/src/account_log.py +++ b/src/account_log.py @@ -1,6 +1,6 @@ """Daily UTF-8 business logs; bounded UI history, no raw requests or message bodies.""" -import hashlib +import json import logging import re from collections import deque @@ -11,20 +11,35 @@ from pathlib import Path UI_LIMIT = 2000 +CREDENTIAL = re.compile( + r"(?i)(cookie|authorization|sessionid|mstoken|a_bogus|password|token|secret|verifyfp|signature|x-bogus|x-tt-params)" +) + + def clean(value): text = str(value).replace("\r", r"\r").replace("\n", r"\n") text = re.sub(r"[\x00-\x1f\x7f\u2028\u2029\u202a-\u202e]", " ", text) - text = re.sub(r"https?://\S+", "[地址已隐藏]", text) text = re.sub( - r"(?i)\b(cookie|authorization|sessionid|msToken|a_bogus|password|token|secret)[\"']?\s*[:=].*", - "[敏感字段已隐藏]", + r"(?i)([?&](?:sessionid|mstoken|a_bogus|token|verifyfp|signature|x-bogus|x-tt-params)=)[^&\s]+", + r"\1[凭据已隐藏]", text, ) - return text[:2000] + text = re.sub( + r"(?i)\b(cookie|authorization)[\"']?\s*[:=].*", + r"\1=[凭据已隐藏]", + text, + ) + return text -def reference(value): - return "N" + hashlib.sha256(str(value).encode()).hexdigest()[:12] +def visible(value, key=""): + if CREDENTIAL.search(str(key)): + return "[凭据已隐藏]" + if isinstance(value, dict): + return {str(k): visible(v, k) for k, v in value.items()} + if isinstance(value, (list, tuple)): + return [visible(item) for item in value] + return value def recent_lines(directory, limit=UI_LIMIT): @@ -86,18 +101,9 @@ class DailyLog(logging.Handler): for key, value in fields.items(): if value is None: continue - if key.lower() in { - "text", - "body", - "cookie", - "headers", - "params", - "result", - "profile", - }: - value = "[已隐藏]" - elif not isinstance(value, (str, int, float, bool)): - value = "[复杂数据未记录]" + value = visible(value, key) + if not isinstance(value, (str, int, float, bool)): + value = json.dumps(value, ensure_ascii=False, default=str) values.append(f"{clean(key)}={clean(value)}") self.logger.log( level, "[%s] %s %s", clean(stage), clean(message), " ".join(values) diff --git a/src/account_session.py b/src/account_session.py index e8bdbc4..dc9658d 100644 --- a/src/account_session.py +++ b/src/account_session.py @@ -9,8 +9,14 @@ from urllib.parse import urlsplit from account_store import decode from douyin_im import EXPRESSION from follow_user import validate_uid -from get_current_user import BROWSER_SCRIPT, compact_user, parse_user_response -from subscribe_notifications import INSTALL, WAIT, detail_request_script +from get_current_user import BROWSER_SCRIPT, business_user, parse_user_response +from subscribe_notifications import ( + INSTALL, + WAIT, + detail_request_script, + history_request_script, + works_request_script, +) class SessionError(RuntimeError): @@ -139,7 +145,7 @@ class Session: async def profile(self): try: - user = compact_user( + user = business_user( parse_user_response( await self.evaluate(BROWSER_SCRIPT, main_world=True) ) @@ -226,6 +232,109 @@ class Session: raise SessionError("通知详情返回了未请求的 ID,已停止处理") return notices + async def history_notices(self): + notices = {} + min_time = max_time = 0 + seen_cursors = set() + for _ in range(1000): + response = await self.json( + history_request_script(min_time, max_time), main_world=True + ) + try: + status = int(response["status"]) + payload = json.loads(response["body"]) + except (KeyError, TypeError, ValueError, json.JSONDecodeError) as exc: + raise SessionError("历史通知响应格式无效") from exc + if status != 200 or payload.get("status_code") != 0: + raise SessionError( + f"历史通知请求失败:HTTP {status},业务码 {payload.get('status_code')}" + ) + rows = payload.get("notice_list_v2") + if rows is None: + rows = [] + if not isinstance(rows, list) or any( + not isinstance(row, dict) for row in rows + ): + raise SessionError("历史通知列表格式无效") + 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) + notices.setdefault(nid, row) + if 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): + raise SessionError("历史通知分页游标未推进") + if any(type(value) not in (int, float) for value in cursor): + raise SessionError("历史通知分页游标格式无效") + seen_cursors.add(cursor) + min_time, max_time = cursor + raise SessionError("历史通知超过 50000 条,已停止以避免无限分页") + + async def works(self, all_pages=False): + profile = await self.identity() + 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 + for _ in range(pages): + response = await self.json( + works_request_script(sec_uid, cursor), main_world=True + ) + try: + status = int(response["status"]) + payload = json.loads(response["body"]) + except (KeyError, TypeError, ValueError, json.JSONDecodeError) as exc: + raise SessionError("作品列表响应格式无效") from exc + if status != 200 or payload.get("status_code") != 0: + raise SessionError( + f"作品列表请求失败:HTTP {status},业务码 {payload.get('status_code')}" + ) + rows = payload.get("aweme_list") + if rows is None: + rows = [] + if not isinstance(rows, list) or any( + not isinstance(row, dict) for row in rows + ): + raise SessionError("作品列表格式无效") + 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) + cover = ((row.get("video") or {}).get("cover") or {}).get( + "url_list" + ) or [] + if not cover and row.get("images"): + cover = (row["images"][0] or {}).get("url_list") or [] + works.setdefault( + ident, + { + "aweme_id": ident, + "desc": row.get("desc") or "", + "create_time": row.get("create_time") or "", + "statistics": row.get("statistics") or {}, + "cover": cover[0] if cover else "", + "business": row, + }, + ) + if not all_pages 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: + raise SessionError("作品列表分页游标格式无效") + if next_cursor in seen_cursors or next_cursor == cursor: + raise SessionError("作品列表分页游标未推进") + seen_cursors.add(next_cursor) + cursor = next_cursor + raise SessionError("作品列表超过 18000 条,已停止以避免无限分页") + async def follow(self, target, check=False): self.require_bound() validate_uid(target) diff --git a/src/account_store.py b/src/account_store.py index b937980..7469156 100644 --- a/src/account_store.py +++ b/src/account_store.py @@ -8,7 +8,7 @@ import uuid from contextlib import contextmanager from pathlib import Path -from account_log import UI_LIMIT, reference +from account_log import UI_LIMIT, visible from follow_user import validate_uid @@ -16,6 +16,8 @@ def validate_rule(rule): rule = dict(rule) rule["require_follow"] = False rule.setdefault("cooldown", 14400) + rule.setdefault("work_mode", "all") + rule.setdefault("work_ids", []) kinds = rule.get("kinds", []) if not isinstance(kinds, list) or any( k not in ("digg", "follow", "comment", "general_notice") for k in kinds @@ -34,6 +36,18 @@ def validate_rule(rule): raise ValueError("执行间隔需为 5..86400 秒") if type(rule["cooldown"]) is not int or not 0 <= rule["cooldown"] <= 31536000: raise ValueError("同 UID 冷却需为 0..31536000 秒") + if rule["work_mode"] not in ("all", "selected"): + raise ValueError("作品监控模式无效") + if ( + not isinstance(rule["work_ids"], list) + or len(rule["work_ids"]) > 50000 + or any( + not isinstance(item, str) or not item.isascii() or not item.isdecimal() + for item in rule["work_ids"] + ) + ): + raise ValueError("作品 ID 列表无效") + rule["work_ids"] = list(dict.fromkeys(rule["work_ids"])) return rule @@ -46,6 +60,8 @@ DEFAULT_RULE = { "text": "", "interval": 30, "cooldown": 14400, + "work_mode": "all", + "work_ids": [], } @@ -56,6 +72,36 @@ def decode(value): raise ValueError("持久化数据损坏,请从备份恢复") from exc +def notice_work_id(notice): + values = [notice] + values.extend( + notice.get(key) + for key in ("comment", "digg", "favorite", "collect", "general_notice") + ) + for value in values: + if not isinstance(value, dict): + continue + for candidate in (value, value.get("aweme"), value.get("item")): + if not isinstance(candidate, dict): + continue + ident = candidate.get("aweme_id") or candidate.get("awemeId") + if isinstance(ident, str) and ident.isascii() and ident.isdecimal(): + return ident + if type(ident) is int and ident > 0: + return str(ident) + return "" + + +def notice_kind(notice): + if notice.get("comment"): + return "comment" + if notice.get("follow"): + return "follow" + if notice.get("digg"): + return "digg" + return "general_notice" + + class Store: def __init__(self, root, audit=None): self.audit = audit @@ -86,12 +132,14 @@ class Store: CREATE TABLE IF NOT EXISTS events ( id INTEGER PRIMARY KEY, source TEXT NOT NULL REFERENCES accounts(id), nid TEXT NOT NULL, created REAL NOT NULL, targets TEXT NOT NULL, rule TEXT NOT NULL, state TEXT NOT NULL, + origin TEXT NOT NULL DEFAULT 'live' CHECK(origin IN ('live','history')), UNIQUE(source,nid)); CREATE TABLE IF NOT EXISTS tasks ( id INTEGER PRIMARY KEY, event INTEGER NOT NULL REFERENCES events(id), source TEXT NOT NULL REFERENCES accounts(id), worker TEXT NOT NULL REFERENCES accounts(id), target TEXT NOT NULL, action TEXT NOT NULL CHECK(action IN ('follow','dm')), params TEXT NOT NULL, dependency INTEGER REFERENCES tasks(id), + priority INTEGER NOT NULL DEFAULT 100, status TEXT NOT NULL DEFAULT 'pending' CHECK(status IN ('pending','running','succeeded','failed','unknown','cancelled')), result TEXT NOT NULL DEFAULT '', updated REAL NOT NULL, UNIQUE(event,worker,target,action)); @@ -109,6 +157,10 @@ class Store: "ALTER TABLE accounts ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0" ) columns = {r["name"] for r in self.db.execute("PRAGMA table_info(inbox)")} + event_columns = { + r["name"] for r in self.db.execute("PRAGMA table_info(events)") + } + task_columns = {r["name"] for r in self.db.execute("PRAGMA table_info(tasks)")} with self.transaction(): if "retry_count" not in columns: self.db.execute( @@ -118,9 +170,20 @@ class Store: self.db.execute( "ALTER TABLE inbox ADD COLUMN next_retry_at REAL NOT NULL DEFAULT 0" ) + if "origin" not in event_columns: + self.db.execute( + "ALTER TABLE events ADD COLUMN origin TEXT NOT NULL DEFAULT 'live' CHECK(origin IN ('live','history'))" + ) + if "priority" not in task_columns: + self.db.execute( + "ALTER TABLE tasks ADD COLUMN priority INTEGER NOT NULL DEFAULT 100" + ) for row in self.db.execute("SELECT id,rule FROM accounts").fetchall(): rule = validate_rule(decode(row["rule"])) - self.db.execute("UPDATE accounts SET rule=? WHERE id=?", (json.dumps(rule), row["id"])) + self.db.execute( + "UPDATE accounts SET rule=? WHERE id=?", + (json.dumps(rule), row["id"]), + ) self.db.executescript(""" CREATE TRIGGER IF NOT EXISTS owner_insert BEFORE INSERT ON accounts WHEN NEW.owner IS NOT NULL AND NOT EXISTS(SELECT 1 FROM accounts WHERE id=NEW.owner AND role='main') @@ -184,10 +247,6 @@ class Store: if row else "不存在的账号" ) - elif key == "nid": - value = reference(value) - elif key == "target": - value = "UID尾号 ***" + (str(value)[-4:] if len(str(value)) > 4 else "") elif key == "action": value = {"follow": "关注", "dm": "私信"}.get(value, value) values[labels.get(key, key)] = value @@ -272,19 +331,10 @@ class Store: def bind_profile(self, ident, profile): uid = validate_uid(profile.get("uid")) - # Store only an explicit allow-list, never cookies, phone numbers or response blobs. - saved = { - key: profile.get(key) - for key in ( - "uid", - "nickname", - "douyin_id", - "following_count", - "follower_count", - "total_favorited", - "aweme_count", - ) - } + saved = visible(profile) + if not isinstance(saved, dict): + raise ValueError("账号资料格式无效") + saved["uid"] = uid with self.transaction(): account = self.account(ident) if account["uid"] is not None and account["uid"] != uid: @@ -298,8 +348,10 @@ class Store: if account["uid"] is None: self.log( "身份绑定", - "首次身份核验通过,绑定已锁定;不记录真实 UID 或平台资料", + "首次身份核验通过,绑定已锁定;业务资料按原文输出,认证凭据保持隐藏", account=ident, + UID=uid, + 平台资料=saved, ) self.db.execute( "UPDATE accounts SET uid=?,nickname=?,avatar=?,profile=? WHERE id=?", @@ -321,7 +373,7 @@ class Store: ) self.log( "规则更新", - "已保存;已入队任务保留原快照,私信正文不写日志", + "已保存;已入队任务保留原快照,业务字段按原文写日志", source=ident, 启用=rule["enabled"], 通知类型=",".join(rule["kinds"]), @@ -330,8 +382,17 @@ class Store: 双动作="关注和私信并行" if rule["follow"] and rule["dm"] else "单动作", 动作间隔秒=rule["interval"], 同组UID冷却秒=rule["cooldown"], + 作品监控模式=rule["work_mode"], + 作品ID=rule["work_ids"], + 私信正文=rule["text"], ) + def set_work_filter(self, ident, mode, ids): + rule = decode(self.account(ident)["rule"]) + rule["work_mode"] = mode + rule["work_ids"] = ids + self.set_rule(ident, rule) + def move(self, worker, owner): with self.transaction(): account = self.account(worker) @@ -460,16 +521,17 @@ class Store: "SELECT COUNT(*) FROM inbox WHERE source=? AND state='pending'", (source,) ).fetchone()[0] - def ingest(self, source, notice): + def ingest(self, source, notice, origin="live"): + if origin not in ("live", "history"): + raise ValueError("事件来源无效") account = self.account(source) 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", ""))) - kind = next( - (k for k in ("digg", "follow", "comment") if notice.get(k)), - "general_notice", + kind = notice_kind(notice) + detail = ( + notice.get(kind) or notice.get("favorite") or notice.get("collect") or {} ) - detail = notice.get(kind) or {} users = detail.get("from_user") or [] if isinstance(users, dict): users = [users] @@ -477,66 +539,100 @@ class Store: users = [detail["comment"]["user"]] targets = sorted({validate_uid(str(u.get("uid", ""))) for u in users}) rule = validate_rule(decode(account["rule"])) - enabled = rule["enabled"] and kind in rule["kinds"] and bool(targets) + work_id = notice_work_id(notice) + work_allowed = ( + kind == "follow" + or rule["work_mode"] == "all" + or work_id in rule["work_ids"] + ) + enabled = ( + rule["enabled"] and kind in rule["kinds"] and bool(targets) and work_allowed + ) + created = notice.get("create_time") or notice.get("createTime") or time.time() + if type(created) not in (int, float) or created <= 0: + created = time.time() + if created > 10_000_000_000: + created /= 1000 with self.transaction(): cur = self.db.execute( - "INSERT OR IGNORE INTO events(source,nid,created,targets,rule,state) VALUES (?,?,?,?,?,?)", + "INSERT OR IGNORE INTO events(source,nid,created,targets,rule,state,origin) VALUES (?,?,?,?,?,?,?)", ( source, nid, - time.time(), + created, json.dumps(targets), json.dumps(rule), "waiting" if enabled else "ignored", + origin, ), ) event = self.db.execute( - "SELECT id FROM events WHERE source=? AND nid=?", (source, nid) - ).fetchone()["id"] + "SELECT id,origin FROM events WHERE source=? AND nid=?", (source, nid) + ).fetchone() + if not cur.rowcount and origin == "live" and event["origin"] == "history": + self.db.execute( + "UPDATE events SET origin='live' WHERE id=?", (event["id"],) + ) + promoted = self.db.execute( + "UPDATE tasks SET priority=100 WHERE event=? AND status='pending'", + (event["id"],), + ).rowcount + self.log( + "新消息优先", + "该通知此前从历史列表进入;已提升未执行任务优先级", + source=source, + event=event["id"], + nid=nid, + 提升任务数=promoted, + ) if cur.rowcount: reason = ( "规则匹配,等待组内分配" if enabled - else ( - "规则未启用,本通知不生成任务" - if not rule["enabled"] - else "通知类型未勾选,本通知不生成任务" - if kind not in rule["kinds"] - else "未提取到有效通知来源用户,不生成任务" - ) + else "规则未启用,本通知不生成任务" + if not rule["enabled"] + else "通知类型未勾选,本通知不生成任务" + if kind not in rule["kinds"] + else "未提取到有效通知来源用户,不生成任务" + if not targets + else "通知不属于当前选择的作品,不生成任务" ) self.log( "详情入库", "详情已取得且所属身份一致", source=source, nid=nid, - event=event, + event=event["id"], + 来源="实时通知" if origin == "live" else "历史列表", 通知类型={ "digg": "点赞", "comment": "评论", "follow": "关注", - "general_notice": "其他", + "general_notice": "其他作品互动", }[kind], + 作品ID=work_id, 目标数量=len(targets), + 业务数据=notice, ) self.log( "规则匹配" if enabled else "通知跳过", reason, source=source, - event=event, + event=event["id"], ) - else: + elif not (origin == "live" and event["origin"] == "history"): self.log( "详情去重", "已有该通知,不重复生成事件或任务", source=source, - event=event, + event=event["id"], nid=nid, ) self.db.execute( "UPDATE inbox SET state='done' WHERE source=? AND nid=?", (source, nid) ) self._dispatch(source) + return bool(cur.rowcount) def dispatch(self, source): with self.transaction(): @@ -545,7 +641,7 @@ class Store: def _dispatch(self, source): main = self.account(source) events = self.db.execute( - "SELECT * FROM events WHERE source=? AND state='waiting' ORDER BY id", + "SELECT * FROM events WHERE source=? AND state='waiting' ORDER BY CASE origin WHEN 'live' THEN 0 ELSE 1 END,id", (source,), ).fetchall() if not events: @@ -600,9 +696,30 @@ class Store: ) continue active = self.db.execute( - "SELECT status FROM tasks WHERE source=? AND target=? AND status IN ('pending','running','unknown') ORDER BY id DESC LIMIT 1", + """SELECT t.id,t.event,t.status,e.origin FROM tasks t JOIN events e ON e.id=t.event + WHERE t.source=? AND t.target=? AND t.status IN ('pending','running','unknown') ORDER BY t.id""", (source, target), - ).fetchone() + ).fetchall() + if event["origin"] == "live": + history = [ + row + for row in active + if row["status"] == "pending" and row["origin"] == "history" + ] + for row in history: + self.db.execute( + "UPDATE tasks SET status='cancelled',result='新消息优先,取消未执行历史任务',updated=? WHERE id=? AND status='pending'", + (time.time(), row["id"]), + ) + self.log( + "新消息优先", + "同 UID 新消息到达,取消尚未执行的历史任务,不重复打扰", + source=source, + event=event["id"], + task=row["id"], + target=target, + ) + active = [row for row in active if row not in history] if active: skipped += 1 self.log( @@ -611,7 +728,7 @@ class Store: source=source, event=event["id"], target=target, - 现有状态=active["status"], + 现有状态=active[-1]["status"], ) continue cooldown = rule.get("cooldown", 14400) @@ -635,7 +752,7 @@ class Store: if not rule[action]: continue cur = self.db.execute( - "INSERT INTO tasks(event,source,worker,target,action,params,dependency,updated) VALUES (?,?,?,?,?,?,NULL,?)", + "INSERT INTO tasks(event,source,worker,target,action,params,dependency,priority,updated) VALUES (?,?,?,?,?,?,NULL,?,?)", ( event["id"], source, @@ -643,6 +760,7 @@ class Store: target, action, json.dumps(rule), + 100 if event["origin"] == "live" else 0, now, ), ) @@ -706,7 +824,7 @@ class Store: first = self.db.execute( """SELECT t.* FROM tasks t WHERE worker=? AND source=? AND status='pending' AND (dependency IS NULL OR EXISTS(SELECT 1 FROM tasks d WHERE d.id=t.dependency AND d.status='succeeded')) - ORDER BY id LIMIT 1""", + ORDER BY priority DESC,id LIMIT 1""", (worker, account["owner"]), ).fetchone() if first is None: diff --git a/src/accounts_app.py b/src/accounts_app.py index f321678..4e22163 100644 --- a/src/accounts_app.py +++ b/src/accounts_app.py @@ -48,6 +48,7 @@ class Backend(QThread): snapshot_ready = Signal(dict) response = Signal(str, bool, str) logs_ready = Signal(list) + data_ready = Signal(str, object) stopped = Signal() def __init__(self, root): @@ -76,7 +77,9 @@ class Backend(QThread): if name == "quit": break try: - await engine.command(name, data) + result = await engine.command(name, data) + if result is not None: + self.data_ready.emit(name, result) self.response.emit(name, True, "操作完成") except Exception as exc: # Only explicitly safe human-facing errors; never emit CDP stack/URLs. @@ -207,6 +210,9 @@ class Window(QMainWindow): ("关闭浏览器", lambda: self.selected_command("close_browser")), ("删除登录数据", self.erase), ("删除小号", self.delete_worker), + ("历史事件操作", self.history_events), + ("获取近期作品", lambda: self.fetch_works(False)), + ("获取全部历史作品", lambda: self.fetch_works(True)), ]: self.button(actions, text, handler) self.tabs.addTab(accounts, "账号组") @@ -269,6 +275,7 @@ class Window(QMainWindow): backend.snapshot_ready.connect(self.refresh) backend.response.connect(self.response) backend.logs_ready.connect(self.append_logs) + backend.data_ready.connect(self.command_data) backend.stopped.connect(self.finished) def button(self, layout, text, handler): @@ -329,6 +336,252 @@ class Window(QMainWindow): ): self.send("delete_worker", {"id": account["id"]}) + def history_events(self): + account = self.selected() + if not account or account["role"] != "main" or not account["uid"]: + QMessageBox.information(self, "提示", "请先选择已登录并绑定身份的大号。") + return + self.send("history_fetch", {"id": account["id"]}) + + def fetch_works(self, all_pages): + 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"]}) + + def command_data(self, name, payload): + if name == "history_fetch": + self.history_dialog(payload) + elif name in ("works_recent", "works_all"): + self.works_dialog(payload) + elif name == "history_enqueue": + QMessageBox.information( + self, + "历史事件已确认", + f"选择 {payload['selected']} 条;新增事件 {payload['events']} 条;创建任务 {payload['tasks']} 条。\n重复、规则不匹配、作品范围外、冷却或自操作事件不会创建任务。新消息优先。", + ) + elif name == "work_filter": + text = ( + "全部作品" + if payload["mode"] == "all" + else f"指定的 {payload['count']} 个作品" + ) + QMessageBox.information(self, "作品监控已保存", "当前监控:" + text) + + @staticmethod + def set_table_checks(table, state): + for row in range(table.rowCount()): + item = table.item(row, 0) + if item is not None: + item.setCheckState(state) + + @staticmethod + def checked_table_data(table): + values = [] + for row in range(table.rowCount()): + item = table.item(row, 0) + if item is not None and item.checkState() == Qt.CheckState.Checked: + values.append(item.data(Qt.ItemDataRole.UserRole)) + return values + + def history_dialog(self, payload): + items = payload["items"] + win = QDialog(self) + 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) + table.setHorizontalHeaderLabels( + [ + "选择", + "时间", + "类型", + "来源UID", + "来源昵称", + "作品ID", + "作品描述", + "评论/内容", + "通知ID", + "完整业务数据", + ] + ) + for row, item in enumerate(items): + check = QTableWidgetItem() + check.setFlags(Qt.ItemFlag.ItemIsEnabled | Qt.ItemFlag.ItemIsUserCheckable) + check.setCheckState(Qt.CheckState.Unchecked) + check.setData(Qt.ItemDataRole.UserRole, item["nid"]) + table.setItem(row, 0, check) + 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"], + item["nid"], + json.dumps(item["business"], ensure_ascii=False, default=str), + ] + for column, value in enumerate(values, 1): + cell = QTableWidgetItem(str(value)) + cell.setToolTip(str(value)) + table.setItem(row, column, cell) + table.horizontalHeader().setSectionResizeMode( + QHeaderView.ResizeMode.ResizeToContents + ) + table.horizontalHeader().setSectionResizeMode(9, QHeaderView.ResizeMode.Stretch) + layout.addWidget(table) + controls = QHBoxLayout() + select_all = QPushButton("全选") + clear = QPushButton("清空选择") + 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 + ) + buttons.accepted.connect(win.accept) + buttons.rejected.connect(win.reject) + layout.addWidget(buttons) + if win.exec() != QDialog.DialogCode.Accepted: + return + selected = self.checked_table_data(table) + if not selected: + QMessageBox.information(self, "没有选择", "未创建任何历史任务。") + return + if ( + QMessageBox.question( + self, + "确认历史批量操作", + f"确认将 {len(selected)} 条历史事件交给小号按当前规则处理?\n这可能产生真实关注和私信;新消息优先,冷却和作品筛选仍会跳过部分事件。", + QMessageBox.StandardButton.Yes | QMessageBox.StandardButton.No, + QMessageBox.StandardButton.No, + ) + == QMessageBox.StandardButton.Yes + ): + self.send("history_enqueue", {"id": payload["account"], "ids": selected}) + + 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) + layout = QVBoxLayout(win) + 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) + table.setHorizontalHeaderLabels( + [ + "选择", + "发布时间", + "作品描述", + "作品ID", + "统计", + "封面URL", + "完整业务数据", + ] + ) + for row, item in enumerate(items): + check = QTableWidgetItem() + check.setFlags(Qt.ItemFlag.ItemIsEnabled | Qt.ItemFlag.ItemIsUserCheckable) + check.setCheckState( + Qt.CheckState.Checked + if item["aweme_id"] in existing + else Qt.CheckState.Unchecked + ) + check.setData(Qt.ItemDataRole.UserRole, item["aweme_id"]) + table.setItem(row, 0, check) + values = [ + item["create_time"], + item["desc"], + item["aweme_id"], + json.dumps(item["statistics"], ensure_ascii=False, default=str), + item["cover"], + json.dumps(item["business"], ensure_ascii=False, default=str), + ] + for column, value in enumerate(values, 1): + cell = QTableWidgetItem(str(value)) + cell.setToolTip(str(value)) + table.setItem(row, column, cell) + table.horizontalHeader().setSectionResizeMode( + QHeaderView.ResizeMode.ResizeToContents + ) + table.horizontalHeader().setSectionResizeMode(6, 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.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.Save + | QDialogButtonBox.StandardButton.Cancel + ) + buttons.accepted.connect(win.accept) + buttons.rejected.connect(win.reject) + layout.addWidget(buttons) + 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)) + if mode == "selected" and not selected: + QMessageBox.warning( + self, + "没有选择作品", + "指定作品模式至少选择一个作品;如需暂停请关闭自动操作规则。", + ) + return + self.send( + "work_filter", {"id": payload["account"], "mode": mode, "ids": selected} + ) + def selected(self): item = self.tree.currentItem() ident = item.data(0, Qt.ItemDataRole.UserRole) if item else None @@ -434,7 +687,7 @@ class Window(QMainWindow): ("digg", "点赞"), ("follow", "关注"), ("comment", "评论"), - ("general_notice", "其他(仅有明确来源 UID 时)"), + ("general_notice", "收藏/其他作品互动(仅有明确来源 UID 时)"), ]: check = QCheckBox(label) check.setChecked(kind in rule["kinds"]) @@ -469,6 +722,8 @@ class Window(QMainWindow): "text": text.toPlainText(), "interval": interval.value(), "cooldown": cooldown.value() * 60, + "work_mode": rule.get("work_mode", "all"), + "work_ids": rule.get("work_ids", []), } try: validate_rule(value) diff --git a/src/get_current_user.py b/src/get_current_user.py index 213ad64..4322a21 100644 --- a/src/get_current_user.py +++ b/src/get_current_user.py @@ -6,7 +6,9 @@ import json import sys from datetime import datetime from pathlib import Path +from typing import Any, cast +from account_log import visible from cdp_explicit import ( # pyright: ignore[reportMissingImports] add_cdp_args, configure, @@ -62,24 +64,14 @@ def parse_user_response(value) -> dict: return {**user, "uid": str(uid)} -def compact_user(user: dict) -> dict: - """保留个人信息常用字段,避免把手机号等敏感字段写入文件。""" - avatars = (user.get("avatar_300x300") or user.get("avatar_thumb") or {}).get( - "url_list" - ) or [] - return { - "uid": user.get("uid"), - "sec_uid": user.get("sec_uid"), - "nickname": user.get("nickname"), - "douyin_id": user.get("short_id") or user.get("unique_id"), - "signature": user.get("signature", ""), - "following_count": user.get("following_count", 0), - "follower_count": user.get("follower_count", 0), - "total_favorited": user.get("total_favorited", 0), - "aweme_count": user.get("aweme_count", 0), - "gender": user.get("gender", 0), - "avatar_url": avatars[0] if avatars else None, - } +def business_user(user: dict) -> dict: + """Return full platform business data while retaining credential protection.""" + result = cast(dict[str, Any], visible(user)) + avatar = result.get("avatar_300x300") or result.get("avatar_thumb") or {} + avatars = avatar.get("url_list") if isinstance(avatar, dict) else [] + result["douyin_id"] = result.get("short_id") or result.get("unique_id") or "" + result["avatar_url"] = avatars[0] if avatars else result.get("avatar_url") or "" + return result def main() -> int: @@ -89,7 +81,7 @@ def main() -> int: args = parser.parse_args() configure(args) try: - user = compact_user(get_user_from_browser()) + user = business_user(get_user_from_browser()) args.output.write_text( json.dumps( {"fetched_at": datetime.now().astimezone().isoformat(), "user": user}, diff --git a/src/subscribe_notifications.py b/src/subscribe_notifications.py index 8e55aa6..67b930e 100644 --- a/src/subscribe_notifications.py +++ b/src/subscribe_notifications.py @@ -185,6 +185,116 @@ def detail_request_script(ids): })(IDS)""".replace("IDS", json.dumps(ids)) +def native_read_script( + endpoint, module_needle, method_needle, export_name, request, expected +): + config = { + "endpoint": endpoint, + "moduleNeedle": module_needle, + "methodNeedle": method_needle, + "exportName": export_name, + "request": request, + "expected": expected, + } + return r"""(async config => { + if (location.origin !== 'https://www.douyin.com') throw Error('抖音页面已关闭'); + const chunks = window.webpackChunkdouyin_web; + if (!chunks) throw Error('未找到抖音运行时'); + let require; + chunks.push([['native-read-' + Date.now()], {}, r => { require = r; }]); + chunks.pop(); + const endpointLiteral = JSON.stringify(config.endpoint); + const entry = Object.entries(require.m).find(([, f]) => { + const source = String(f); + return source.includes(config.moduleNeedle) && source.includes(endpointLiteral) + && (!config.methodNeedle || source.includes(config.methodNeedle)); + }); + const api = entry && require(entry[0]); + const method = config.exportName ? api?.[config.exportName] : Object.values(api || {}).find( + value => typeof value === 'function' && String(value).includes(endpointLiteral) + && (!config.methodNeedle || String(value).includes(config.methodNeedle))); + const client = window.axiosInstance; + if (typeof method !== 'function' || !client?.interceptors?.response) + throw Error('只读请求 SDK 未就绪或已变化'); + let raw, timer; + const observer = client.interceptors.response.use(response => { + const params = response.config?.params || {}; + const url = new URL(response.config?.url || '', location.origin); + const matched = Object.entries(config.expected).every(([key, value]) => String(params[key]) === String(value)); + const xhr = response.request; + if (url.origin === location.origin && url.pathname === config.endpoint && matched + && xhr && (!xhr.responseType || xhr.responseType === 'text') + && typeof xhr.responseText === 'string') + raw = {status: response.status, body: xhr.responseText}; + return response; + }); + try { + await Promise.race([ + method(config.request), + new Promise((_, reject) => { timer = setTimeout(() => reject(Error('timeout')), 20000); }) + ]); + } catch (_) { + throw Error('只读 SDK 请求失败或超时,请检查登录与网络'); + } finally { + clearTimeout(timer); + client.interceptors.response.eject(observer); + } + if (!raw) throw Error('无法读取只读请求原始响应,SDK 可能已变化'); + return JSON.stringify(raw); + })(CONFIG)""".replace("CONFIG", json.dumps(config, ensure_ascii=False)) + + +def history_request_script(min_time=0, max_time=0, count=50): + if type(min_time) not in (int, float) or type(max_time) not in (int, float): + raise ValueError("历史通知游标无效") + if type(count) is not int or not 1 <= count <= 50: + raise ValueError("历史通知每页数量需为 1..50") + params = { + "notice_group": 700, + "count": count, + "min_time": min_time, + "max_time": max_time, + "is_mark_read": 0, + } + return native_read_script( + "/aweme/v1/web/notice/", + "/aweme/v1/web/notice/", + "is_new_notice", + "", + params, + params, + ) + + +def works_request_script(sec_uid, cursor=0, count=18): + if ( + not isinstance(sec_uid, str) + or not sec_uid + or len(sec_uid) > 256 + or any(ord(char) < 32 for char in sec_uid) + ): + raise ValueError("作品列表身份参数无效") + if type(cursor) is not int or cursor < 0: + raise ValueError("作品列表游标无效") + if type(count) is not int or not 1 <= count <= 50: + raise ValueError("作品列表每页数量需为 1..50") + request = { + "userId": sec_uid, + "maxCursor": cursor, + "count": count, + "needTimeList": True, + } + expected = {"sec_user_id": sec_uid, "max_cursor": cursor, "count": count} + return native_read_script( + "/aweme/v1/web/aweme/post/", + "/aweme/v1/web/aweme/post/", + "", + "fetchUserPost", + request, + expected, + ) + + def details(ids, uid): expression = detail_request_script(ids) # 推送可能先于详情入库;仅对这次事件重试,不做定时通知扫描。 diff --git a/src/test_accounts.py b/src/test_accounts.py index ed01404..f0e7252 100644 --- a/src/test_accounts.py +++ b/src/test_accounts.py @@ -39,6 +39,17 @@ def notice(uid, nid="111", target="999"): return {"user_id": uid, "nid_str": nid, "follow": {"from_user": [{"uid": target}]}} +def work_notice(uid, nid, target, work, kind="digg"): + return { + "user_id": uid, + "nid_str": nid, + kind: { + "from_user": [{"uid": target, "nickname": "完整昵称"}], + "aweme": {"aweme_id": work, "desc": "完整作品描述"}, + }, + } + + class StoreTest(unittest.TestCase): def setUp(self): self.temp = tempfile.TemporaryDirectory() @@ -63,7 +74,9 @@ class StoreTest(unittest.TestCase): validate_rule(rule(interval=True)) with self.assertRaises(ValueError): validate_rule(rule(cooldown=True)) - normalized = validate_rule({k: v for k, v in rule(require_follow=True).items() if k != "cooldown"}) + normalized = validate_rule( + {k: v for k, v in rule(require_follow=True).items() if k != "cooldown"} + ) self.assertFalse(normalized["require_follow"]) self.assertEqual(normalized["cooldown"], 14400) self.assertFalse(json.loads(self.store.account(self.other)["rule"])["enabled"]) @@ -91,7 +104,9 @@ class StoreTest(unittest.TestCase): self.assertEqual(saved["uid"], "301") self.assertEqual(saved["nickname"], "自动昵称") self.assertEqual(json.loads(saved["profile"])["follower_count"], 123) - self.assertNotIn("must-not-store", saved["profile"]) + stored_profile = json.loads(saved["profile"]) + self.assertEqual(stored_profile["cookie"], "[凭据已隐藏]") + self.assertEqual(stored_profile["phone"], "must-not-store") with self.assertRaises(ValueError): self.store.bind_profile(b, profile) self.assertIsNone(self.store.account(b)["uid"]) @@ -141,7 +156,12 @@ class StoreTest(unittest.TestCase): def test_dedup_round_robin_and_parallel_batch(self): second = self.store.add("worker", "第二小号", "202", self.main) - for nid, target in (("111", "901"), ("111", "901"), ("112", "902"), ("113", "903")): + for nid, target in ( + ("111", "901"), + ("111", "901"), + ("112", "902"), + ("113", "903"), + ): self.store.ingest(self.main, notice("101", nid, target)) tasks = sorted(self.store.tasks(), key=lambda t: t["id"]) self.assertEqual(len(tasks), 6) @@ -150,13 +170,85 @@ class StoreTest(unittest.TestCase): ) batch = self.store.claim_batch(self.worker) self.assertEqual({task["action"] for task in batch}, {"follow", "dm"}) - self.assertEqual({task["status"] for task in self.store.tasks() if task["id"] in {t["id"] for t in batch}}, {"running"}) + self.assertEqual( + { + task["status"] + for task in self.store.tasks() + if task["id"] in {t["id"] for t in batch} + }, + {"running"}, + ) self.assertIsNone(self.store.claim(self.worker)) self.store.finish(batch[0]["id"], "failed", {}) self.store.finish(batch[1]["id"], "unknown", {}) next_batch = self.store.claim_batch(self.worker) self.assertEqual({task["target"] for task in next_batch}, {"903"}) + def test_selected_work_filter_covers_work_interactions_but_not_follow(self): + self.store.set_rule( + self.main, + rule( + dm=False, + cooldown=0, + kinds=["digg", "follow", "comment", "general_notice"], + work_mode="selected", + work_ids=["700"], + ), + ) + self.assertTrue( + self.store.ingest(self.main, work_notice("101", "1011", "901", "700")) + ) + self.assertTrue( + self.store.ingest(self.main, work_notice("101", "1012", "902", "701")) + ) + self.assertTrue( + self.store.ingest( + self.main, work_notice("101", "1013", "903", "700", "favorite") + ) + ) + self.assertTrue(self.store.ingest(self.main, notice("101", "1014", "904"))) + self.assertEqual(len(self.store.tasks()), 3) + states = dict( + self.store.db.execute( + "SELECT nid,state FROM events WHERE source=?", (self.main,) + ) + ) + self.assertEqual(states["1012"], "ignored") + self.assertEqual( + {states[nid] for nid in ("1011", "1013", "1014")}, {"dispatched"} + ) + + 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") + self.store.ingest(self.main, notice("101", "2002", "902"), origin="history") + self.store.ingest(self.main, notice("101", "2003", "903"), origin="live") + first = self.store.claim_batch(self.worker) + self.assertEqual({task["target"] for task in first}, {"903"}) + for task in first: + self.store.finish(task["id"], "succeeded", {}) + self.store.ingest(self.main, notice("101", "2004", "901"), origin="live") + old = self.store.db.execute( + "SELECT status FROM tasks WHERE event=(SELECT id FROM events WHERE nid='2001')" + ).fetchall() + self.assertEqual({row["status"] for row in old}, {"cancelled"}) + promoted = notice("101", "2005", "904") + self.store.ingest(self.main, promoted, origin="history") + self.store.ingest(self.main, promoted, origin="live") + event = self.store.db.execute( + "SELECT id,origin FROM events WHERE source=? AND nid='2005'", (self.main,) + ).fetchone() + self.assertEqual(event["origin"], "live") + self.assertEqual( + { + row["priority"] + for row in self.store.db.execute( + "SELECT priority FROM tasks WHERE event=?", (event["id"],) + ) + }, + {100}, + ) + def test_group_cooldown_persists_is_shared_and_other_groups_are_independent(self): second = self.store.add("worker", "第二小号", "202", self.main) other_worker = self.store.add("worker", "其他组小号", "203", self.other) @@ -172,18 +264,26 @@ class StoreTest(unittest.TestCase): batch = self.store.claim_batch(self.worker) self.assertEqual(len(batch), 2) self.assertEqual( - self.store.db.execute("SELECT last_at FROM cooldowns WHERE source=? AND target='999'", (self.main,)).fetchone()[0], 100 + self.store.db.execute( + "SELECT last_at FROM cooldowns WHERE source=? AND target='999'", + (self.main,), + ).fetchone()[0], + 100, ) for task in batch: self.store.finish(task["id"], "succeeded", {}) # The same UID in a different main-account group is independent. self.store.ingest(self.other, notice("102", "211", "999")) - self.assertEqual(len([t for t in self.store.tasks() if t["worker"] == other_worker]), 2) + self.assertEqual( + len([t for t in self.store.tasks() if t["worker"] == other_worker]), 2 + ) self.store.close() self.store = Store(self.temp.name) with patch("account_store.time.time", return_value=200): self.store.ingest(self.main, notice("101", "113", "999")) - self.assertEqual(len([t for t in self.store.tasks() if t["source"] == self.main]), 2) + self.assertEqual( + len([t for t in self.store.tasks() if t["source"] == self.main]), 2 + ) with patch("account_store.time.time", return_value=14501): self.store.ingest(self.main, notice("101", "114", "999")) new_tasks = [t for t in self.store.tasks() if t["source"] == self.main] @@ -239,7 +339,9 @@ class StoreTest(unittest.TestCase): self.store.finish(batch[0]["id"], "failed", {}) self.store.finish(batch[1]["id"], "unknown", {}) self.assertEqual(self.store.claim_batch(self.worker), []) - self.assertEqual({t["status"] for t in self.store.tasks()}, {"failed", "unknown"}) + self.assertEqual( + {t["status"] for t in self.store.tasks()}, {"failed", "unknown"} + ) def test_disabled_and_cross_group_identity(self): with self.assertRaises(ValueError): @@ -297,7 +399,9 @@ class StoreTest(unittest.TestCase): self.assertEqual(migrated["cooldown"], 14400) self.assertFalse(migrated["require_follow"]) self.assertIsNotNone( - self.store.db.execute("SELECT 1 FROM sqlite_master WHERE type='table' AND name='cooldowns'").fetchone() + self.store.db.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' AND name='cooldowns'" + ).fetchone() ) def test_legacy_inbox_migration_preserves_pending_ids(self): @@ -553,16 +657,18 @@ class EngineTest(unittest.IsolatedAsyncioTestCase): await self.engine.start([ident]) self.assertEqual(self.engine.desired, set()) - async def test_im_result_classification_and_redaction(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": "secret"}, + "message": {"client_id": "c1", "server_id": "s1", "content": "完整私信正文"}, + "token": "AUTH-CREDENTIAL", } status, result = await self.engine.perform(session, task) self.assertEqual(status, "succeeded") - self.assertNotIn("secret", json.dumps(result)) + self.assertIn("完整私信正文", json.dumps(result, ensure_ascii=False)) + self.assertNotIn("AUTH-CREDENTIAL", json.dumps(result)) session.im.return_value = {"error": "SDK_REQUEST_FAILED"} self.assertEqual((await self.engine.perform(session, task))[0], "unknown") session.im.return_value = {"error": "LOGIN_REQUIRED"} @@ -663,6 +769,125 @@ class EngineTest(unittest.IsolatedAsyncioTestCase): self.assertEqual(len(self.engine.store.tasks()), 2) self.assertIn("不阻塞新通知", self.engine.states[self.main]) + async def test_history_and_works_paginate_read_only_with_exact_ids(self): + session = Session(None, "http://127.0.0.1:9222", "101") + first = { + "status": 200, + "body": json.dumps( + { + "status_code": 0, + "notice_list_v2": [notice("101", "9223372036854775701")], + "has_more": 1, + "min_time": 100, + "max_time": 200, + } + ), + } + second = { + "status": 200, + "body": json.dumps( + { + "status_code": 0, + "notice_list_v2": [notice("101", "9223372036854775702")], + "has_more": 0, + "min_time": 90, + "max_time": 190, + } + ), + } + session.json = AsyncMock(side_effect=[first, second]) + history = await session.history_notices() + self.assertEqual( + [row["nid_str"] for row in history], + ["9223372036854775701", "9223372036854775702"], + ) + self.assertTrue( + all( + "is_mark_read" in call.args[0] and "700" in call.args[0] + for call in session.json.await_args_list + ) + ) + session.identity = AsyncMock(return_value={"uid": "101", "sec_uid": "SEC-UID"}) + session.json = AsyncMock( + side_effect=[ + { + "status": 200, + "body": json.dumps( + { + "status_code": 0, + "aweme_list": [ + { + "aweme_id": "9223372036854775703", + "desc": "完整作品", + "create_time": 123, + "author": {"uid": "101"}, + "statistics": {"digg_count": 9}, + "video": { + "cover": { + "url_list": [ + "https://example.invalid/full.jpg" + ] + } + }, + } + ], + "has_more": 1, + "max_cursor": 321, + } + ), + }, + { + "status": 200, + "body": json.dumps( + { + "status_code": 0, + "aweme_list": [], + "has_more": 0, + "max_cursor": 0, + } + ), + }, + ] + ) + works = await session.works(True) + self.assertEqual(works[0]["aweme_id"], "9223372036854775703") + self.assertEqual(works[0]["desc"], "完整作品") + self.assertEqual(works[0]["cover"], "https://example.invalid/full.jpg") + self.assertEqual(session.json.await_count, 2) + + async def test_history_and_works_reject_cursor_loops_and_foreign_identity(self): + session = Session(None, "http://127.0.0.1:9222", "101") + repeating = { + "status": 200, + "body": json.dumps( + { + "status_code": 0, + "notice_list_v2": [], + "has_more": 1, + "min_time": 0, + "max_time": 0, + } + ), + } + session.json = AsyncMock(return_value=repeating) + with self.assertRaisesRegex(SessionError, "游标未推进"): + await session.history_notices() + session.identity = AsyncMock(return_value={"uid": "101", "sec_uid": "SEC-UID"}) + session.json = AsyncMock( + return_value={ + "status": 200, + "body": json.dumps( + { + "status_code": 0, + "aweme_list": [{"aweme_id": "700", "author": {"uid": "other"}}], + "has_more": 0, + } + ), + } + ) + with self.assertRaisesRegex(SessionError, "身份不符"): + await session.works(False) + async def test_detail_64bit_ids_preserved(self): session = Session(None, "http://127.0.0.1:9222", "101") session.identity = AsyncMock(return_value={"uid": "101"}) @@ -681,6 +906,52 @@ class EngineTest(unittest.IsolatedAsyncioTestCase): result = await session.details([nid]) self.assertEqual(str(result[0]["nid"]), nid) + async def test_history_preview_confirmation_and_work_filter_commands(self): + session = AsyncMock() + history = [work_notice("101", "3001", "901", "700")] + works = [ + { + "aweme_id": "700", + "desc": "完整作品描述", + "create_time": 123, + "statistics": {"digg_count": 5}, + "cover": "https://example.invalid/cover.jpg", + "business": {"aweme_id": "700", "desc": "完整作品描述"}, + } + ] + session.history_notices.return_value = history + session.works.return_value = works + self.engine.store.set_rule( + self.main, + rule(kinds=["digg", "follow", "comment", "general_notice"]), + ) + self.engine.open = AsyncMock(return_value=session) + preview = await self.engine.command("history_fetch", {"id": self.main}) + self.assertEqual(preview["items"][0]["nid"], "3001") + self.assertEqual(preview["items"][0]["work_id"], "700") + self.assertEqual(preview["items"][0]["actor_names"], ["完整昵称"]) + self.assertEqual(self.engine.store.tasks(), []) + with self.assertRaises(ValueError): + await self.engine.command( + "history_enqueue", {"id": self.main, "ids": ["not-loaded"]} + ) + result = await self.engine.command( + "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"]) + saved = await self.engine.command( + "work_filter", {"id": self.main, "mode": "selected", "ids": ["700"]} + ) + 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() + async def test_move_during_identity_check_does_not_cross_groups(self): other = self.engine.store.add("main", "另一组", "102") self.engine.store.set_rule(other, rule()) @@ -737,10 +1008,14 @@ class EngineTest(unittest.IsolatedAsyncioTestCase): if entered == 2: both_entered.set() await release.wait() - return {"status": "succeeded"} if kind == "follow" else { - "success": True, - "message": {"client_id": "client", "server_id": "server"}, - } + return ( + {"status": "succeeded"} + if kind == "follow" + else { + "success": True, + "message": {"client_id": "client", "server_id": "server"}, + } + ) async def follow_action(*args, **kwargs): return await action("follow", *args) diff --git a/src/test_activity.py b/src/test_activity.py index df65b98..b9ae998 100644 --- a/src/test_activity.py +++ b/src/test_activity.py @@ -72,14 +72,10 @@ def test_sensitive_data_single_line_and_file_failure(): ) audit.record("安全", 'authorization="AUTH_SECRET"') text = all_logs(audit) + assert "PRIVATE_BODY" in text and "https://site.invalid/" in text assert not any( secret in text - for secret in ( - "PRIVATE_BODY", - "COOKIE_SECRET", - "AUTH_SECRET", - "URL_SECRET", - ) + for secret in ("COOKIE_SECRET", "AUTH_SECRET", "URL_SECRET") ) assert "\u2029" not in text and len(text.splitlines()) == 2 finally: @@ -124,7 +120,7 @@ def test_business_stages_skips_dependencies_and_rollback(): assert not store.tasks() assert "触发者就是执行小号自身" in all_logs(audit) assert "生成任务数=0" in all_logs(audit) - assert "1234567890123456789" not in all_logs(audit) + assert "1234567890123456789" in all_logs(audit) store.record_push(main, ["2"]) store.ingest(main, notification("2")) batch = store.claim_batch(worker) @@ -152,7 +148,7 @@ def test_business_stages_skips_dependencies_and_rollback(): "人工核对", ): assert f"[{stage}]" in text - assert "PRIVATE_TEMPLATE" not in text + assert "PRIVATE_TEMPLATE" in text with pytest.raises(ValueError), store.transaction(): store.log("测试", "ROLLBACK_MUST_NOT_APPEAR") raise ValueError("rollback") diff --git a/src/test_browser_matrix.py b/src/test_browser_matrix.py index b5a34f8..bb4b0ba 100644 --- a/src/test_browser_matrix.py +++ b/src/test_browser_matrix.py @@ -36,9 +36,69 @@ def fixture_for(expected): "nickname": "offline fixture", "follower_count": 42, "aweme_count": 3, + "sec_uid": "SEC-" + expected, }, } ) + elif "/aweme/v1/web/notice/?" in route.request.url: + from urllib.parse import parse_qs, urlsplit + + query = parse_qs(urlsplit(route.request.url).query) + assert query["is_mark_read"] == ["0"] and query["notice_group"] == ["700"] + first = query["min_time"] == ["0"] and query["max_time"] == ["0"] + await route.fulfill( + json={ + "status_code": 0, + "notice_list_v2": [ + { + "nid": 9223372036854775701 + if first + else 9223372036854775702, + "nid_str": "9223372036854775701" + if first + else "9223372036854775702", + "user_id": expected, + "create_time": 100 if first else 90, + "digg": { + "from_user": [{"uid": "999", "nickname": "完整昵称"}], + "aweme": {"aweme_id": "800", "desc": "完整作品描述"}, + }, + } + ], + "has_more": 1 if first else 0, + "min_time": 100 if first else 90, + "max_time": 200 if first else 190, + } + ) + elif "/aweme/v1/web/aweme/post/?" in route.request.url: + from urllib.parse import parse_qs, urlsplit + + query = parse_qs(urlsplit(route.request.url).query) + assert query["sec_user_id"] == ["SEC-" + expected] + first = query["max_cursor"] == ["0"] + await route.fulfill( + json={ + "status_code": 0, + "aweme_list": [ + { + "aweme_id": "9223372036854775703", + "desc": "完整作品描述", + "create_time": 123, + "author": {"uid": expected}, + "statistics": {"digg_count": 9}, + "video": { + "cover": { + "url_list": ["https://example.invalid/full.jpg"] + } + }, + } + ] + if first + else [], + "has_more": 1 if first else 0, + "max_cursor": 321 if first else 0, + } + ) elif "/aweme/v1/web/notice/detail/" in route.request.url: from urllib.parse import parse_qs, urlsplit @@ -73,16 +133,18 @@ SDK_FIXTURE = r"""uid => { class NoticeFrontier {} NoticeFrontier.frontierInstance=f; const observers=new Map([[0,response=>response]]);let token=0; window.axiosInstance={interceptors:{response:{use:fn=>{observers.set(++token,fn);return token;},eject:id=>observers.delete(id)}}}; - const api={getNoticeDetail:params=>new Promise((resolve,reject)=>{ - const path='/aweme/v1/web/notice/detail/',xhr=new XMLHttpRequest(); - xhr.open('GET',path+'?'+new URLSearchParams({...params,sdk_fixture:'1'})); + const request=(path,params)=>new Promise((resolve,reject)=>{ + const xhr=new XMLHttpRequest();xhr.open('GET',path+'?'+new URLSearchParams(params)); xhr.onload=()=>{const response={status:xhr.status,config:{url:path,params},request:xhr}; for(const observe of observers.values())observe(response); resolve(JSON.parse(xhr.responseText));}; xhr.onerror=()=>reject(Error('fixture XHR failed'));xhr.send(); - })}; - const req=id=>id==='notice'?{NoticeFrontier}:id==='api'?api:{decodedFrame:()=>({service:20313,payload:new TextEncoder().encode('[]')})}; - req.m={notice:function(){/* NOTICE_PUSH_EVENT_NAMES:function */},codec:function(){/* .decodedFrame= .encodeFrame= */},api:function(){/* getNoticeDetail: */}}; + }); + const api={getNoticeDetail:params=>request('/aweme/v1/web/notice/detail/',{...params,sdk_fixture:'1'})}; + function history(params){"/aweme/v1/web/notice/";"is_new_notice";return request('/aweme/v1/web/notice/',params);} + function fetchUserPost(params){"/aweme/v1/web/aweme/post/";return request('/aweme/v1/web/aweme/post/',{sec_user_id:params.userId,max_cursor:params.maxCursor,count:params.count});} + const req=id=>id==='notice'?{NoticeFrontier}:id==='api'?api:id==='history'?{x:history}:id==='posts'?{fetchUserPost}:{decodedFrame:()=>({service:20313,payload:new TextEncoder().encode('[]')})}; + req.m={notice:function(){/* NOTICE_PUSH_EVENT_NAMES:function */},codec:function(){/* .decodedFrame= .encodeFrame= */},api:function(){/* getNoticeDetail: */},history:function(){return "/aweme/v1/web/notice/"+"is_new_notice";},posts:function(){return "/aweme/v1/web/aweme/post/";}}; window.webpackChunkdouyin_web=[]; window.webpackChunkdouyin_web.push=function(chunk){Array.prototype.push.call(this,chunk);chunk[2](req);}; window.__offlineSdk={observers:()=>observers.size,count:()=>listeners.get('message').size, @@ -193,6 +255,23 @@ async def matrix(chrome): results[ "native_detail_transport_partial_null_precision_and_cleanup" ] = True + history = await sessions[0].history_notices() + assert [row["nid_str"] for row in history] == [ + "9223372036854775701", + "9223372036854775702", + ] + works = await sessions[0].works(True) + assert len(works) == 1 and works[0]["aweme_id"] == "9223372036854775703" + assert works[0]["desc"] == "完整作品描述" + assert ( + await sessions[0].evaluate( + "window.__offlineSdk.observers()", main_world=True + ) + == 1 + ) + results[ + "native_history_and_works_pagination_identity_precision_and_cleanup" + ] = True results["identity_isolated"] = ( await sessions[0].evaluate("localStorage.identity") == "101" and await sessions[1].evaluate("localStorage.identity") == "102" diff --git a/src/test_onboarding_ui.py b/src/test_onboarding_ui.py index 8b164d9..d906fa9 100644 --- a/src/test_onboarding_ui.py +++ b/src/test_onboarding_ui.py @@ -9,7 +9,15 @@ from unittest.mock import patch os.environ.setdefault("QT_QPA_PLATFORM", "offscreen") -from PySide6.QtWidgets import QApplication, QDialog, QLineEdit, QMessageBox +from PySide6.QtCore import Qt +from PySide6.QtWidgets import ( + QApplication, + QCheckBox, + QDialog, + QLineEdit, + QMessageBox, + QTableWidget, +) from account_store import DEFAULT_RULE from accounts_app import Backend, Window @@ -135,6 +143,104 @@ class OnboardingUiTest(unittest.TestCase): self.window.delete_worker() send.assert_not_called() + def test_history_and_works_buttons_require_bound_main(self): + main = {"id": "main", "role": "main", "uid": "101"} + with ( + patch.object(self.window, "selected", return_value=main), + patch.object(self.window, "send") as send, + ): + self.window.history_events() + self.window.fetch_works(False) + self.window.fetch_works(True) + self.assertEqual( + [call.args for call in send.call_args_list], + [ + ("history_fetch", {"id": "main"}), + ("works_recent", {"id": "main"}), + ("works_all", {"id": "main"}), + ], + ) + + def test_history_preview_requires_selection_and_explicit_confirmation(self): + payload = { + "account": "main", + "items": [ + { + "nid": "9007199254740993", + "create_time": 123, + "kind": "comment", + "work_id": "700", + "work_desc": "完整作品描述", + "actor_uids": ["999"], + "actor_names": ["完整昵称"], + "comment": "完整评论正文", + "business": { + "nid_str": "9007199254740993", + "comment": "完整评论正文", + }, + } + ], + } + + def choose(dialog): + table = dialog.findChild(QTableWidget) + assert ( + table is not None and table.item(0, 9).text().find("完整评论正文") >= 0 + ) + table.item(0, 0).setCheckState(Qt.CheckState.Checked) + return QDialog.DialogCode.Accepted + + with ( + patch.object(QDialog, "exec", choose), + patch.object( + QMessageBox, "question", return_value=QMessageBox.StandardButton.Yes + ), + patch.object(self.window, "send") as send, + ): + self.window.history_dialog(payload) + send.assert_called_once_with( + "history_enqueue", {"id": "main", "ids": ["9007199254740993"]} + ) + + def test_works_preview_saves_selected_mode_and_full_business_data(self): + payload = { + "account": "main", + "mode": "selected", + "selected": [], + "all_pages": True, + "items": [ + { + "aweme_id": "700", + "desc": "完整作品描述", + "create_time": 123, + "statistics": {"collect_count": 8}, + "cover": "https://example.invalid/full.jpg", + "business": {"aweme_id": "700", "desc": "完整作品描述"}, + } + ], + } + + def choose(dialog): + table = dialog.findChild(QTableWidget) + assert table is not None and "完整作品描述" in table.item(0, 6).text() + table.item(0, 0).setCheckState(Qt.CheckState.Checked) + all_mode = next( + box + for box in dialog.findChildren(QCheckBox) + if "监控全部作品" in box.text() + ) + all_mode.setChecked(False) + return QDialog.DialogCode.Accepted + + with ( + patch.object(QDialog, "exec", choose), + patch.object(self.window, "send") as send, + ): + self.window.works_dialog(payload) + send.assert_called_once_with( + "work_filter", {"id": "main", "mode": "selected", "ids": ["700"]} + ) + def test_anonymous_profile_is_never_bound(self): for user in ({}, {"uid": "0"}, {"uid": True}, {"uid": "not-a-uid"}): with self.assertRaises(RuntimeError): diff --git a/src/test_subscribe_notifications.py b/src/test_subscribe_notifications.py index 6d98bda..0777902 100644 --- a/src/test_subscribe_notifications.py +++ b/src/test_subscribe_notifications.py @@ -6,6 +6,7 @@ import subprocess from contextlib import redirect_stdout from pathlib import Path from tempfile import TemporaryDirectory +from typing import Any, cast from unittest.mock import patch import subscribe_notifications as sub @@ -212,6 +213,64 @@ window.webpackChunkdouyin_web.push=function(chunk){chunk[2](req);return Array.pr assert json.loads(raw["body"])["notice_list_v2"][0]["nid"] == 9007199254740993 +def test_native_history_and_works_transport(): + history_script = cast(Any, sub.history_request_script) + for args in ((0, 0, 0), (0, 0, 51), ("bad", 0, 10)): + try: + history_script(*args) + except ValueError: + pass + else: + raise AssertionError("必须拒绝无效历史分页参数") + for args in (("", 0, 18), ("SEC", -1, 18), ("SEC", 0, 51)): + try: + sub.works_request_script(*args) + except ValueError: + pass + else: + raise AssertionError("必须拒绝无效作品分页参数") + script = r""" +const assert=require('node:assert/strict'); +const input=JSON.parse(require('node:fs').readFileSync(0,'utf8')); +const observers=new Map([[0,response=>response]]);let seq=0; +const client={interceptors:{response:{use:fn=>{observers.set(++seq,fn);return seq;},eject:id=>observers.delete(id)}}}; +const respond=(endpoint,params,body)=>{ + if(input.missing)return Promise.resolve({}); + const response={status:200,config:{url:endpoint,params},request:{responseType:'text',responseText:body}}; + for(const callback of observers.values())assert.equal(callback(response),response); + return Promise.resolve(JSON.parse(body)); +}; +function notice(params){"/aweme/v1/web/notice/";"is_new_notice";assert.equal(params.is_mark_read,0);assert.equal(params.notice_group,700);return respond('/aweme/v1/web/notice/',params,'{"status_code":0,"notice_list_v2":[],"has_more":0}');} +function fetchUserPost(params){"/aweme/v1/web/aweme/post/";const query={sec_user_id:params.userId,max_cursor:params.maxCursor,count:params.count};return respond('/aweme/v1/web/aweme/post/',query,'{"status_code":0,"aweme_list":[],"has_more":0}');} +const modules={notice:{x:notice},posts:{fetchUserPost}}; +const req=id=>modules[id]; +req.m={notice:function(){return "/aweme/v1/web/notice/"+"is_new_notice";},posts:function(){return "/aweme/v1/web/aweme/post/";}}; +global.location={origin:'https://www.douyin.com'}; +global.window={axiosInstance:client,webpackChunkdouyin_web:[]}; +window.webpackChunkdouyin_web.push=function(chunk){chunk[2](req);return Array.prototype.push.call(this,chunk);}; +(async()=>{ + if(input.missing)await assert.rejects(()=>eval(input.expression)); + else {const result=JSON.parse(await eval(input.expression));assert.equal(result.status,200);} + assert.deepEqual([...observers.keys()],[0]); + console.log('ok'); +})().catch(error=>{console.error(error);process.exitCode=1;}); +""" + expressions = [ + sub.history_request_script(0, 0, 50), + sub.works_request_script("SEC-UID", 0, 18), + ] + for expression in expressions: + for missing in (False, True): + subprocess.run( + ["node", "-e", script], + input=json.dumps({"expression": expression, "missing": missing}), + text=True, + capture_output=True, + check=True, + timeout=10, + ) + + def test_js_bridge(): # 假 Webpack/FWS;只测自有监听器,不向真实账号发送通知。 script = r""" @@ -273,5 +332,6 @@ if __name__ == "__main__": test_details() test_subscription_and_cli() test_native_detail_transport() + test_native_history_and_works_transport() test_js_bridge() print("subscription checks passed")