diff --git a/WINDOWS-README.txt b/WINDOWS-README.txt index 7251c24..779bc45 100644 --- a/WINDOWS-README.txt +++ b/WINDOWS-README.txt @@ -1,7 +1,7 @@ -抖音账号助手 0.1.5 — Windows 10/11 x64 +抖音账号助手 0.1.6 — Windows 10/11 x64 【安装】 -只需运行 DouyinAccounts-0.1.5-Windows-x64-Setup.exe。 +只需运行 DouyinAccounts-0.1.6-Windows-x64-Setup.exe。 无需安装 Python、Qt、Patchright、Node.js 或另行下载浏览器。 安装程序不会自动启动业务;没有开机自启、自动登录或默认私信。 默认使用随包的 fingerprint-chromium(adryfish 项目)和 Patchright 驱动,不再回退到系统 Chrome。 @@ -16,14 +16,16 @@ 账号已添加或登录身份与原绑定不一致时会提示,不会覆盖已有绑定。登录成功也不会自动启动业务。 若关闭了浏览器,可点击“登录 / 刷新信息”继续;已绑定账号点击此按钮可刷新资料和统计。 只保留一个抖音首页或自己的页面;多个匹配标签会拒绝执行,防止选错。 -3. 选择大号,配置规则。默认关闭;可配置通知类别、关注、私信、正文和顺序。 - 勾选大号后点击启动;规则未开启不会产生或执行业务任务。 +3. 选择大号,配置规则。默认关闭;可配置通知类别、关注、私信、正文、不同任务批次间隔和同组 UID 冷却。 + 冷却默认 240 分钟(4 小时),设为 0 可关闭;勾选大号后点击启动。规则未开启不会产生或执行业务任务。 4. 每条匹配通知只给一个关联小号,按组内轮询,不跨组操作。 一条聚合通知可能包含多个来源用户,每个有效来源用户各执行配置动作。 5. 小号不可用时等待原小号,不自动改派;大号不可用时整个组暂停。 6. 停止只暂停业务、等待在途请求保存,浏览器保留。 关闭浏览器为独立操作,先停止账号组;退出软件也不会自动退出 Chrome。 7. “待核对”表示执行结果不确定,不会重发。请到平台人工核对,再确认成功或失败。 + 同一通知同时开启关注和私信时,两个动作并行开始、结果互不依赖,不等待默认 30 秒批次间隔。 + 批次间隔仍用于同一小号的不同通知批次;旧版“必须关注成功后再私信”配置升级后自动关闭。 8. “正在监听通知(N 条详情待重试,不阻塞新通知)”表示部分通知暂未返回详情。 这些 ID 会保留,按 30、60、120、240、300 秒间隔重试,之后最多每 5 分钟一次;不拖停整组,不影响正常新通知。 身份不符或真正的接口错误仍会暂停;不能根据空详情推断通知一定已删除。 @@ -35,10 +37,19 @@ 此操作不删除账号定义、任务和历史,也不影响其他账号。 更换 Chrome 可执行路径会使用不同登录目录,不复制 Cookie。 +【同组 UID 冷却与双动作并行(0.1.6 新增)】 +冷却按大号组共享:该组任一小号开始操作某 UID 后,组内其他小号在设定时长内也不会再次操作该 UID;不同大号组互不影响。 +冷却从任务实际领取、即将请求平台时开始并持久化,程序重启不重置。只排队后因停止、切换归属或删除小号而取消的任务不会启动冷却。 +冷却期间,上一次结果无论成功、明确失败还是待核对,都不会因新通知再次创建任务;待核对记录在人工处理前还会继续阻止重复任务,优先保证不重复打扰。 +在任务实际开始前,如果同组同 UID 已存在待执行、执行中或待核对任务,新通知也会直接跳过,避免短时间多条通知重复排队。 +修改冷却时长会用于之后收到的通知;现有冷却按新通知当时的规则计算剩余时间。日志记录冷却开始、剩余秒数、已有任务跳过和所属本地大号名称,不记录完整 UID。 +同一通知同时启用关注和私信时,两项组成一个并行批次,同时调用平台接口;其中一项失败或待核对不会取消另一项。两项各自保存结果。 +默认 30 秒“不同任务批次间隔”只作用于前后两个通知批次,不插入到同一并行批次的关注和私信之间。 + 【详细运行记录(0.1.5 新增)】 “运行记录”显示最近 2000 条,重新打开程序也会加载最近的日志。可关闭自动滚动,或打开日志文件夹。 完整记录按电脑本地日期写入 %LOCALAPPDATA%\DouyinAccounts\logs\YYYY-MM-DD.log,UTF-8 编码,不因 UI 上限而截断,也不自动删除旧日期文件。 -记录推送接收与去重、详情请求/重试、规则不匹配、无执行小号、自操作保护、轮询分配、任务入队/执行/结果、依赖等待、人工核对、账号操作和异常。 +记录推送接收与去重、详情请求/重试、规则不匹配、无执行小号、自操作保护、UID 冷却、并行批次、轮询分配、任务入队/执行/结果、旧任务依赖等待、人工核对、账号操作和异常。 例如“触发者就是执行小号自身”会明确显示跳过原因;“分配完成”不代表已经关注或发送私信。 监听无新通知时每分钟写一次心跳,重复等待状态不反复刷屏。任务列表也仅展示最近 2000 条,完整历史仍在数据库中。 日志不保存 Cookie、签名、私信正文、原始请求/响应;通知标识使用稳定摘要,目标 UID 只显示尾号,账号用本地名称标识。 @@ -59,7 +70,8 @@ 卸载/覆盖升级默认保留账号数据。从 0.1.0 升级旧数据库时会先备份再迁移,保留原 UID、归属、任务和历史。 升级 0.1.2 前请先停止业务并关闭旧 Chrome。新内核使用独立登录目录,需要自行重新登录原账号。 旧目录及旧浏览器配置保留,但新版本不自动沿用旧 Chrome 的路径配置;如需自定义请重新设置指纹浏览器路径。 -0.1.3 及更新版本首次启动旧库时会增加通知重试字段;0.1.5 增加删除标记字段,旧账号默认均保留,旧通知 ID、任务和登录目录不变。 +0.1.3 及更新版本首次启动旧库时会增加通知重试字段;0.1.5 增加删除标记字段;0.1.6 增加组内 UID 冷却表和规则字段。旧账号、通知 ID、任务和登录目录不变。 +0.1.6 将旧规则的“必须关注成功后再私信”自动关闭;旧版已排队且带前置依赖的任务仍按原快照安全处理,不重写历史。 0.1.4 修复了新点赞、评论推送后无法取得详情的问题:复用网页自身请求层,自动补齐当前运行时参数,保持原始 ID 精度。 本次已经实际验证“新增点赞和评论 → 两条推送 → 两条完整详情”,未执行自动关注或私信。 备份前先停止业务、退出软件并关闭所有账号浏览器,备份整个数据目录。 diff --git a/build_windows.py b/build_windows.py index 4549132..87d2a5b 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.5", + "app_version": "0.1.6", "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.5-Windows-x64-Setup.exe" + installer = ROOT / "release/DouyinAccounts-0.1.6-Windows-x64-Setup.exe" (installer.parent / "SHA256SUMS.txt").write_text( sha256(installer) + " " + installer.name + "\n", encoding="ascii" ) diff --git a/docs/2026-09-07-00-01-group-uid-cooldown-and-parallel-actions.md b/docs/2026-09-07-00-01-group-uid-cooldown-and-parallel-actions.md new file mode 100644 index 0000000..68739b4 --- /dev/null +++ b/docs/2026-09-07-00-01-group-uid-cooldown-and-parallel-actions.md @@ -0,0 +1,216 @@ +# 0.1.6:大号组共享 UID 冷却与关注/私信并行 + +## 需求与已确认口径 + +用户要求: + +1. 同一个 UID 可配置冷却时长,默认 4 小时,期间不重复打扰。 +2. 同一通知同时配置关注和私信时,两项并行执行,不等待默认 30 秒动作间隔。 +3. 开始本轮修改前先推送已有版本。 + +已先将完整 0.1.5 快照提交并推送至 `origin/main`: + +```text +fee3951 feat: add Windows multi-account assistant +``` + +通过交互确认的业务口径: + +- 冷却按**大号组共享**:一个大号组内所有执行小号共用;不同大号组互不影响。 +- 同时开启关注和私信时**强制并行**。旧“私信前必须关注成功”不再控制新任务;两个结果分别保存,互不取消。 + +## 冷却生命周期 + +### 为什么不在通知入库时立即开始 + +通知可能刚入队就因用户停止业务、切换小号归属或删除小号而取消。如果通知入库即开始 4 小时冷却,会在没有实际请求平台时消耗冷却期。 + +因此冷却从任务被执行小号**实际领取、即将发起平台请求**时开始。领取批次和写入冷却记录处于同一个 SQLite 事务中;进程重启后仍然有效。 + +### 冷却开始前的重复保护 + +同一大号组、同一目标 UID 如果已经存在以下任一状态的任务,新通知不再创建重复任务: + +- `pending`:已排队; +- `running`:正在执行; +- `unknown`:结果不确定,必须人工核对。 + +这样可以关闭“第一条还没执行,短时间又来了多条点赞/评论”的重复排队窗口。只排队后被取消且从未领取的任务不会写入冷却表。 + +`unknown` 比普通冷却更保守:即使配置冷却时间已经过去,只要结果仍未人工核对,仍不自动创建同 UID 的新任务。这延续了项目原有的“结果不确定绝不自动重发”约束。 + +### 冷却期行为 + +任务领取后,不论最终状态为: + +- 成功; +- 明确失败; +- 结果不确定; + +该大号组在配置时长内都不会因为新通知再次创建同 UID 的任务。日志会写出“冷却开始”或“冷却跳过”、本地大号名称及剩余秒数,但不写完整 UID。 + +默认值: + +```text +14400 秒 = 240 分钟 = 4 小时 +``` + +规则 UI 使用分钟配置: + +- 范围 0–525600 分钟; +- 0 表示关闭; +- 内部继续使用整数秒存储; +- 数据层允许 0–31536000 秒。 + +新通知使用其入库时的规则快照判断冷却时长。如果用户修改规则,之后收到的通知使用新时长;历史事件和任务不被重写。 + +### 持久化结构 + +新增表: + +```sql +CREATE TABLE cooldowns ( + 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) +); +``` + +主键 `(source, target)` 使冷却天然按大号组和目标 UID 共享。执行小号不在主键中,因此组内轮询切换到另一个小号不会绕过冷却。另一个大号组使用不同 `source`,不会互相影响。 + +表只保留每组每 UID 最近一次实际领取时间,采用覆盖更新,不为同一组合无限追加历史行。详细历史仍由每日运行日志和任务表保存。 + +## 关注和私信并行 + +### 任务创建 + +同一事件、同一目标、同一执行小号同时开启关注和私信时,仍创建两个独立任务,但两者: + +- `dependency=NULL`; +- 保存同一规则快照; +- 属于同一事件和目标; +- 不再创建“私信依赖关注任务”的关系。 + +UI 删除了“必须关注成功后才发送私信”复选框。兼容字段 `require_follow` 暂时保留在 JSON 中,但规则验证、保存和旧库升级都会强制归一化为 `false`,避免旧界面配置继续产生新依赖。 + +升级不会重写历史任务。0.1.5 之前已经入队且带 dependency 的旧任务,领取逻辑仍保留兼容处理:按旧快照等待或取消,不将旧私信突然并行发送。 + +### 批次领取 + +新增 `Store.claim_batch(worker)`: + +1. 使用 `BEGIN IMMEDIATE` 串行化领取事务; +2. 拒绝账号未绑定、归属不符、规则关闭或已有 running 任务; +3. 选出最早可执行任务; +4. 若同一 worker/source/event/target 下还有无依赖的 pending 兄弟任务,一并改为 running; +5. 在同一事务写入共享冷却时间; +6. 返回一个或两个任务。 + +旧数据库中的 `worker_single_running` 部分唯一索引会删除,因为合法并行批次需要同一小号同时有两个 running 任务。重复领取仍由 `BEGIN IMMEDIATE`、running 检查和状态更新保证;另一个 Store 连接在前一个事务提交后只能看到 running,不会再次领取。 + +保留 `Store.claim()` 供单任务测试和兼容调用;产品 Worker 使用 `claim_batch()`。 + +### 同时执行和独立结果 + +`Engine.perform_batch()` 使用: + +```python +await asyncio.gather(..., return_exceptions=True) +``` + +关注和私信协程同时启动。单项抛出异常时: + +- 该项记为 `unknown / INTERRUPTED_CHECK_MANUALLY`; +- 另一项不会被取消,继续完成; +- 两项分别调用 `Store.finish()` 保存状态和耗时; +- unknown 仍需用户人工核对,不自动重发。 + +默认 30 秒设置改名为“不同任务批次间隔”。Worker 只在领取下一个通知批次前检查它;同一批次内部不插入 sleep,也不等待关注结果。 + +停止业务仍等待在途并行批次完成,再退出循环和保留浏览器,避免关闭过程中把两项静默丢失。程序意外中断时,启动恢复会把仍为 running 的两项都改为 unknown。 + +## 可观察日志 + +0.1.5 的按日日志新增: + +- `[冷却开始]`:任务已实际领取、组共享冷却开始、冷却秒数; +- `[冷却跳过]`:仍在冷却期及剩余秒数; +- `[冷却跳过]`:已有 pending/running/unknown 时不重复排队; +- `[并行批次]`:同一通知的关注和私信同时开始; +- 两条独立 `[任务开始]` 和两条独立 `[任务结果]`。 + +日志继续不保存 Cookie、签名、私信正文、原始响应、完整通知 ID 或完整目标 UID。UI 仍只展示最近 2000 条,完整日志按本地日期保留在: + +```text +%LOCALAPPDATA%\DouyinAccounts\logs\YYYY-MM-DD.log +``` + +## 数据库升级 + +0.1.6 启动时: + +1. `CREATE TABLE IF NOT EXISTS cooldowns`; +2. `DROP INDEX IF EXISTS worker_single_running`; +3. 所有账号规则缺少 `cooldown` 时补为 14400; +4. 旧 `require_follow` 统一改为 false。 + +旧账号、UID 绑定、归属、通知、任务、日志和浏览器目录均保留。旧任务的 params 快照不批量改写。 + +## 测试 + +### Python 回归 + +Linux 和 Windows 完整 pytest 均: + +```text +64 passed +``` + +新增或调整的核心验证: + +- 缺少 cooldown 的旧规则升级为 14400 秒; +- 旧 require_follow 升级为 false; +- cooldowns 表自动创建; +- bool 等非法类型被拒绝,0 可关闭冷却; +- 通知入队但未领取时不写冷却; +- 同 UID 已有 pending 任务时,新通知不重复创建; +- 领取关注/私信批次后只写一条组共享冷却; +- 同组第二个小号不能绕过冷却; +- 不同大号组对同 UID 互不影响; +- SQLite 重开后冷却仍有效; +- 4 小时到期后新通知可再次创建任务; +- cooldown=0 时前一批完成后新事件可再次创建; +- 双动作批次同时变成 running; +- 一个动作 failed、另一个 unknown 时结果独立,不生成自动重试; +- 程序中断后两个 running 任务均恢复为 unknown; +- Worker 测试使用闸门:关注和私信必须都进入协程后才放行,证明不是串行; +- 点击停止时,在途并行批次完成后才返回;关注和私信各调用一次且均落库成功。 + +### Windows 浏览器与打包 + +- Windows pytest:64 项通过。 +- fingerprint-chromium/Patchright 离线矩阵全部通过:端口/目录隔离、固定指纹种子、SDK 详情请求、64 位通知 ID、登录回填、重连、错误身份拒绝、小号删除和兄弟账号隔离等均保持。 +- 矩阵 `real_write_actions=0`,未使用真实账号执行关注或私信。 +- 打包后 GUI smoke:Qt、核心模块、SQLite、Patchright 和内置浏览器路径均通过;临时账号数 0、真实写动作 0。 +- 6 个相关 Python 文件 LSP error 检查为 0。 +- 本地 23 个 Python 源文件 SHA256 与 Windows `build-manifest.json` 全部一致。 + +本次没有使用真实账号验证并行写操作,避免产生未授权关注或私信。并发调度通过可控制的异步闸门测试验证,关注与私信各自的 SDK 适配继续由原有独立测试覆盖。 + +## 发布 + +```text +C:\Users\rogee\Desktop\抖音账号助手-0.1.6\ + DouyinAccounts-0.1.6-Windows-x64-Setup.exe + SHA256SUMS.txt + 使用说明.txt +``` + +- 安装包大小:222,651,315 字节。 +- SHA256:`6ab893fc2c7de8dbad4dcca43339eb2a2919d340f7b3ef6fe6b4a84cf2c038bc`。 +- 桌面副本已重新计算校验和。 +- 浏览器与依赖版本不变:fingerprint-chromium 148.0.7778.215、Patchright 1.62.3、Python 3.12.10。 + +未代用户安装或启动真实业务。覆盖安装后,建议先打开大号规则确认“同组同 UID 冷却”为 240 分钟,再手动启动账号组。旧版已排队任务仍按原历史快照处理;新通知使用并行与冷却规则。 diff --git a/installer.iss b/installer.iss index 2a07870..8bdf5a3 100644 --- a/installer.iss +++ b/installer.iss @@ -2,14 +2,14 @@ [Setup] AppId={{C9C3EE3B-3666-4E41-A52B-BB309838F701} AppName=抖音账号助手 -AppVersion=0.1.5 +AppVersion=0.1.6 DefaultDirName={localappdata}\Programs\DouyinAccounts DefaultGroupName=抖音账号助手 PrivilegesRequired=lowest ArchitecturesAllowed=x64compatible ArchitecturesInstallIn64BitMode=x64compatible OutputDir=release -OutputBaseFilename=DouyinAccounts-0.1.5-Windows-x64-Setup +OutputBaseFilename=DouyinAccounts-0.1.6-Windows-x64-Setup Compression=lzma2 SolidCompression=yes WizardStyle=modern diff --git a/src/account_engine.py b/src/account_engine.py index 44c8bf1..b7e3b46 100644 --- a/src/account_engine.py +++ b/src/account_engine.py @@ -311,6 +311,17 @@ class Engine: return "failed", {"code": result["status_code"]} return "unknown", {"code": "SDK_RESULT_UNCONFIRMED"} + async def perform_batch(self, session, tasks): + results = await asyncio.gather( + *(self.perform(session, task) for task in tasks), return_exceptions=True + ) + return [ + ("unknown", {"code": "INTERRUPTED_CHECK_MANUALLY"}) + if isinstance(result, BaseException) + else result + for result in results + ] + async def worker_loop(self, ident): session = None while ident in self.desired: @@ -349,23 +360,17 @@ class Engine: self.state(ident, f"等待动作间隔({interval} 秒)") await self.pause(ident, 1) continue - task = self.store.claim(ident) - if task: - self.state(ident, "正在执行任务") - started = time.monotonic() - try: - status, result = await self.perform(session, task) - except Exception: - status, result = ( - "unknown", - {"code": "INTERRUPTED_CHECK_MANUALLY"}, - ) - self.store.finish( - task["id"], - status, - result, - round((time.monotonic() - started) * 1000), + tasks = self.store.claim_batch(ident) + if tasks: + self.state( + ident, + "正在并行执行关注和私信" if len(tasks) > 1 else "正在执行任务", ) + started = time.monotonic() + results = await self.perform_batch(session, tasks) + elapsed = round((time.monotonic() - started) * 1000) + for task, (status, result) in zip(tasks, results, strict=True): + self.store.finish(task["id"], status, result, elapsed) self.last_write[ident] = time.monotonic() else: self.state(ident, "等待组内任务") diff --git a/src/account_store.py b/src/account_store.py index 143bbfa..b937980 100644 --- a/src/account_store.py +++ b/src/account_store.py @@ -13,6 +13,9 @@ from follow_user import validate_uid def validate_rule(rule): + rule = dict(rule) + rule["require_follow"] = False + rule.setdefault("cooldown", 14400) kinds = rule.get("kinds", []) if not isinstance(kinds, list) or any( k not in ("digg", "follow", "comment", "general_notice") for k in kinds @@ -27,10 +30,10 @@ def validate_rule(rule): raise ValueError("启用规则必须选择通知类型及动作") if rule["dm"] and not rule["text"].strip(): raise ValueError("私信正文不能为空") - if rule["require_follow"] and not (rule["follow"] and rule["dm"]): - raise ValueError("先关注成功再私信必须同时启用两个动作") if type(rule.get("interval")) is not int or not 5 <= rule["interval"] <= 86400: raise ValueError("执行间隔需为 5..86400 秒") + if type(rule["cooldown"]) is not int or not 0 <= rule["cooldown"] <= 31536000: + raise ValueError("同 UID 冷却需为 0..31536000 秒") return rule @@ -42,6 +45,7 @@ DEFAULT_RULE = { "require_follow": False, "text": "", "interval": 30, + "cooldown": 14400, } @@ -91,7 +95,11 @@ class Store: 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)); - CREATE UNIQUE INDEX IF NOT EXISTS worker_single_running ON tasks(worker) WHERE status='running'; + CREATE TABLE IF NOT EXISTS cooldowns ( + 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)); + DROP INDEX IF EXISTS worker_single_running; """) self.migrate_accounts() if "deleted" not in { @@ -110,6 +118,9 @@ class Store: self.db.execute( "ALTER TABLE inbox ADD COLUMN next_retry_at REAL NOT NULL DEFAULT 0" ) + 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.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') @@ -304,7 +315,7 @@ class Store: def set_rule(self, ident, rule): if self.account(ident)["role"] != "main": raise ValueError("只有大号可配置规则") - validate_rule(rule) + rule = validate_rule(rule) self.db.execute( "UPDATE accounts SET rule=? WHERE id=?", (json.dumps(rule), ident) ) @@ -316,8 +327,9 @@ class Store: 通知类型=",".join(rule["kinds"]), 关注=rule["follow"], 私信=rule["dm"], - 先关注=rule["require_follow"], + 双动作="关注和私信并行" if rule["follow"] and rule["dm"] else "单动作", 动作间隔秒=rule["interval"], + 同组UID冷却秒=rule["cooldown"], ) def move(self, worker, owner): @@ -587,37 +599,64 @@ class Store: target=target, ) continue - dependency = None + active = self.db.execute( + "SELECT status FROM tasks WHERE source=? AND target=? AND status IN ('pending','running','unknown') ORDER BY id DESC LIMIT 1", + (source, target), + ).fetchone() + if active: + skipped += 1 + self.log( + "冷却跳过", + "同组同 UID 已有待执行、执行中或待核对任务,不重复排队", + source=source, + event=event["id"], + target=target, + 现有状态=active["status"], + ) + continue + cooldown = rule.get("cooldown", 14400) + previous = self.db.execute( + "SELECT last_at FROM cooldowns WHERE source=? AND target=?", + (source, target), + ).fetchone() + now = time.time() + if cooldown and previous and previous["last_at"] + cooldown > now: + skipped += 1 + self.log( + "冷却跳过", + "同组同 UID 仍在冷却期,本通知不创建任务", + source=source, + event=event["id"], + target=target, + 剩余秒=max(1, round(previous["last_at"] + cooldown - now)), + ) + continue for action in ("follow", "dm"): if not rule[action]: continue cur = self.db.execute( - "INSERT INTO tasks(event,source,worker,target,action,params,dependency,updated) VALUES (?,?,?,?,?,?,?,?)", + "INSERT INTO tasks(event,source,worker,target,action,params,dependency,updated) VALUES (?,?,?,?,?,?,NULL,?)", ( event["id"], source, worker["id"], target, action, - event["rule"], - dependency if rule["require_follow"] else None, - time.time(), + json.dumps(rule), + now, ), ) created += 1 self.log( "任务入队", - "已创建任务,等待执行", + "已创建任务,等待执行;同一通知的关注和私信将并行", source=source, worker=worker["id"], event=event["id"], task=cur.lastrowid, target=target, action=action, - 前置任务=dependency if rule["require_follow"] else None, ) - if action == "follow": - dependency = cur.lastrowid self.log( "分配完成", "分配流程已完成;不代表动作已执行", @@ -632,7 +671,7 @@ class Store: ) self.db.execute("UPDATE accounts SET rr=? WHERE id=?", (rr, source)) - def claim(self, worker): + def _claim(self, worker, batch): with self.transaction(): account = self.account(worker) if ( @@ -640,14 +679,14 @@ class Store: or account["owner"] is None or account["uid"] is None ): - return None + return [] owner = self.account(account["owner"]) if owner["uid"] is None or not decode(owner["rule"])["enabled"]: - return None + return [] if self.db.execute( "SELECT 1 FROM tasks WHERE worker=? AND status='running'", (worker,) ).fetchone(): - return None + return [] cancelled = self.db.execute( """UPDATE tasks SET status='cancelled',result='前置关注未成功',updated=? WHERE worker=? AND status='pending' AND dependency IN @@ -657,20 +696,20 @@ class Store: for row in cancelled: self.log( "任务取消", - "前置关注未成功,取消依赖步骤,不重发", + "前置关注未成功,取消旧版依赖步骤,不重发", worker=worker, source=owner["id"], task=row["id"], event=row["event"], action=row["action"], ) - task = self.db.execute( + 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""", (worker, account["owner"]), ).fetchone() - if task is None: + if first is None: blocked = self.db.execute( "SELECT t.id,t.dependency,d.status FROM tasks t JOIN tasks d ON d.id=t.dependency WHERE t.worker=? AND t.source=? AND t.status='pending' ORDER BY t.id LIMIT 1", (worker, account["owner"]), @@ -679,7 +718,7 @@ class Store: if key and self._dependency_wait.get(worker) != key: self.log( "依赖等待", - "前置关注尚未确认成功;待核对不会自动重发", + "旧版前置关注尚未确认成功;待核对不会自动重发", worker=worker, source=owner["id"], task=blocked["id"], @@ -687,23 +726,65 @@ class Store: 前置状态=blocked["status"], ) self._dependency_wait[worker] = key - return None + return [] self._dependency_wait.pop(worker, None) - self.db.execute( + tasks = [first] + if batch and first["dependency"] is None: + tasks = self.db.execute( + """SELECT * FROM tasks WHERE worker=? AND source=? AND event=? AND target=? + AND status='pending' AND dependency IS NULL ORDER BY id""", + (worker, first["source"], first["event"], first["target"]), + ).fetchall() + now = time.time() + self.db.executemany( "UPDATE tasks SET status='running',updated=? WHERE id=?", - (time.time(), task["id"]), + [(now, task["id"]) for task in tasks], ) - self.log( - "任务开始", - "已领取任务,即将发起动作;中断后需人工核对,不自动重发", - source=task["source"], - worker=worker, - event=task["event"], - task=task["id"], - target=task["target"], - action=task["action"], - ) - return dict(task) + cooldown = decode(first["params"]).get("cooldown", 14400) + if cooldown: + self.db.execute( + """INSERT INTO cooldowns(source,target,last_at,event) VALUES (?,?,?,?) + ON CONFLICT(source,target) DO UPDATE SET last_at=excluded.last_at,event=excluded.event""", + (first["source"], first["target"], now, first["event"]), + ) + self.log( + "冷却开始", + "任务已实际领取;从现在起组内所有小号共享该 UID 冷却", + source=first["source"], + worker=worker, + event=first["event"], + target=first["target"], + 冷却秒=cooldown, + ) + if len(tasks) > 1: + self.log( + "并行批次", + "同一通知的关注和私信同时开始,互不依赖且不等待动作间隔", + source=first["source"], + worker=worker, + event=first["event"], + target=first["target"], + 任务数=len(tasks), + ) + for task in tasks: + self.log( + "任务开始", + "已领取任务,即将发起动作;中断后需人工核对,不自动重发", + source=task["source"], + worker=worker, + event=task["event"], + task=task["id"], + target=task["target"], + action=task["action"], + ) + return [dict(task) for task in tasks] + + def claim(self, worker): + tasks = self._claim(worker, False) + return tasks[0] if tasks else None + + def claim_batch(self, worker): + return self._claim(worker, True) def finish(self, task_id, status, result, elapsed_ms=None): if status not in ("succeeded", "failed", "unknown"): diff --git a/src/accounts_app.py b/src/accounts_app.py index 344f2f5..f321678 100644 --- a/src/accounts_app.py +++ b/src/accounts_app.py @@ -421,12 +421,11 @@ class Window(QMainWindow): return rule = decode(account["rule"]) win, form, buttons = dialog("自动操作规则 — " + account["name"], self) - enabled, follow, dm, require = [QCheckBox() for _ in range(4)] + enabled, follow, dm = [QCheckBox() for _ in range(3)] for widget, key, label in [ (enabled, "enabled", "启用自动操作"), (follow, "follow", "执行关注"), (dm, "dm", "发送私信"), - (require, "require_follow", "必须关注成功后才发送私信"), ]: widget.setChecked(rule[key]) form.addRow(label, widget) @@ -447,9 +446,15 @@ class Window(QMainWindow): interval = QSpinBox() interval.setRange(5, 86400) interval.setValue(rule["interval"]) - form.addRow("小号动作间隔(秒)", interval) + form.addRow("不同任务批次间隔(秒)", interval) + cooldown = QSpinBox() + cooldown.setRange(0, 525600) + cooldown.setSpecialValueText("关闭") + cooldown.setSuffix(" 分钟") + cooldown.setValue(rule.get("cooldown", 14400) // 60) + form.addRow("同组同 UID 冷却", cooldown) warning = QLabel( - "每条通知只分配给一个小号,组内轮询。已排队任务保留原规则快照;关闭规则会暂停执行。\n启用后只处理新通知,结果不确定不会重发;请确保行为获得授权并符合平台规则。" + "每条通知只分配给一个小号,组内轮询。同一 UID 的冷却由整个大号组共享,默认 240 分钟;冷却从任务实际开始时计算。\n关注和私信同时开启时会并行执行、结果互不依赖,不等待批次间隔。已排队任务保留原规则快照;结果不确定不会重发。" ) warning.setWordWrap(True) form.addRow(warning) @@ -459,10 +464,11 @@ class Window(QMainWindow): "enabled": enabled.isChecked(), "follow": follow.isChecked(), "dm": dm.isChecked(), - "require_follow": require.isChecked(), + "require_follow": False, "kinds": [k for k, v in kinds.items() if v.isChecked()], "text": text.toPlainText(), "interval": interval.value(), + "cooldown": cooldown.value() * 60, } try: validate_rule(value) diff --git a/src/test_accounts.py b/src/test_accounts.py index 5178088..ed01404 100644 --- a/src/test_accounts.py +++ b/src/test_accounts.py @@ -29,7 +29,7 @@ def rule(**overrides): "enabled": True, "follow": True, "dm": True, - "require_follow": True, + "require_follow": False, "text": "offline-only-test", **overrides, } @@ -61,6 +61,11 @@ class StoreTest(unittest.TestCase): validate_rule(rule(text="")) with self.assertRaises(ValueError): 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"}) + self.assertFalse(normalized["require_follow"]) + self.assertEqual(normalized["cooldown"], 14400) self.assertFalse(json.loads(self.store.account(self.other)["rule"])["enabled"]) with self.assertRaises(ValueError): self.store.add("worker", "坏UID", "01") @@ -134,23 +139,68 @@ class StoreTest(unittest.TestCase): with self.assertRaises(sqlite3.IntegrityError): self.store.add("worker", "不能归属小号", owner=self.worker) - def test_dedup_round_robin_and_dependency(self): + def test_dedup_round_robin_and_parallel_batch(self): second = self.store.add("worker", "第二小号", "202", self.main) - for nid in ("111", "111", "112", "113"): - self.store.ingest(self.main, notice("101", nid)) + 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) self.assertEqual( [tasks[i]["worker"] for i in (0, 2, 4)], [self.worker, second, self.worker] ) - first = self.store.claim(self.worker) - assert first is not None - self.assertEqual(first["action"], "follow") + 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.assertIsNone(self.store.claim(self.worker)) - self.store.finish(first["id"], "succeeded", {}) - second_task = self.store.claim(self.worker) - assert second_task is not None - self.assertEqual(second_task["action"], "dm") + 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_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) + self.store.set_rule(self.other, rule()) + with patch("account_store.time.time", return_value=100): + self.store.ingest(self.main, notice("101", "111", "999")) + self.assertEqual( + self.store.db.execute("SELECT count(*) FROM cooldowns").fetchone()[0], 0 + ) + # A second notification before execution cannot duplicate the queued bundle. + self.store.ingest(self.main, notice("101", "112", "999")) + self.assertEqual(len(self.store.tasks()), 2) + 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 + ) + 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.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) + 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] + self.assertEqual(len(new_tasks), 4) + self.assertEqual({t["worker"] for t in new_tasks[:2]}, {second}) + + def test_zero_cooldown_allows_later_event_after_first_bundle_finishes(self): + self.store.set_rule(self.main, rule(cooldown=0)) + self.store.ingest(self.main, notice("101", "111", "999")) + batch = self.store.claim_batch(self.worker) + for task in batch: + self.store.finish(task["id"], "succeeded", {}) + self.store.ingest(self.main, notice("101", "112", "999")) + self.assertEqual(len(self.store.tasks()), 4) + self.assertEqual( + self.store.db.execute("SELECT count(*) FROM cooldowns").fetchone()[0], 0 + ) def test_move_cancels_only_pending_and_is_atomic(self): self.store.ingest(self.main, notice("101")) @@ -169,28 +219,27 @@ class StoreTest(unittest.TestCase): self.store.move(self.worker, None) self.assertIsNone(self.store.account(self.worker)["owner"]) - def test_restart_running_becomes_unknown_never_resends(self): + def test_restart_running_batch_becomes_unknown_never_resends(self): self.store.ingest(self.main, notice("101")) - first = self.store.claim(self.worker) - assert first is not None + batch = self.store.claim_batch(self.worker) + self.assertEqual(len(batch), 2) self.store.close() self.store = Store(self.temp.name) self.store.recover() - self.assertIsNone(self.store.claim(self.worker)) - self.store.resolve(first["id"], True) - second_task = self.store.claim(self.worker) - assert second_task is not None - self.assertEqual(second_task["action"], "dm") + self.assertEqual({t["status"] for t in self.store.tasks()}, {"unknown"}) + self.assertEqual(self.store.claim_batch(self.worker), []) + for task in batch: + self.store.resolve(task["id"], True) + self.assertEqual(self.store.claim_batch(self.worker), []) - def test_unknown_failed_and_no_auto_retry(self): + def test_parallel_results_are_independent_and_never_resend(self): self.store.ingest(self.main, notice("101")) - first = self.store.claim(self.worker) - assert first is not None - self.store.finish(first["id"], "failed", {}) - self.assertIsNone(self.store.claim(self.worker)) - self.assertEqual( - {t["status"] for t in self.store.tasks()}, {"failed", "cancelled"} - ) + batch = self.store.claim_batch(self.worker) + self.assertEqual(len(batch), 2) + 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"}) def test_disabled_and_cross_group_identity(self): with self.assertRaises(ValueError): @@ -236,6 +285,21 @@ class StoreTest(unittest.TestCase): ).fetchone() self.assertEqual(tuple(row), ("pending", 20, 400)) + def test_rule_migration_adds_default_cooldown_and_disables_old_dependency(self): + legacy = rule(require_follow=True) + legacy.pop("cooldown") + self.store.db.execute( + "UPDATE accounts SET rule=? WHERE id=?", (json.dumps(legacy), self.main) + ) + self.store.close() + self.store = Store(self.temp.name) + migrated = json.loads(self.store.account(self.main)["rule"]) + 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() + ) + def test_legacy_inbox_migration_preserves_pending_ids(self): self.store.record_push(self.main, ["111"]) self.store.db.executescript(""" @@ -663,26 +727,41 @@ class EngineTest(unittest.IsolatedAsyncioTestCase): worker.page.is_closed = lambda: False self.engine.sessions[self.main] = main self.engine.open = AsyncMock(return_value=worker) - entered = asyncio.Event() + both_entered = asyncio.Event() release = asyncio.Event() + entered = 0 - async def action(*args): - entered.set() + async def action(kind, *args): + nonlocal entered + entered += 1 + if entered == 2: + both_entered.set() await release.wait() - return {"status": "succeeded"} + return {"status": "succeeded"} if kind == "follow" else { + "success": True, + "message": {"client_id": "client", "server_id": "server"}, + } - worker.follow.side_effect = action + async def follow_action(*args, **kwargs): + return await action("follow", *args) + + async def dm_action(*args, **kwargs): + return await action("dm", *args) + + worker.follow.side_effect = follow_action + worker.im.side_effect = dm_action runner = asyncio.create_task(self.engine.worker_loop(self.worker)) self.engine.runners[self.worker] = runner - await asyncio.wait_for(entered.wait(), 2) + await asyncio.wait_for(both_entered.wait(), 2) stopping = asyncio.create_task(self.engine.stop()) await asyncio.sleep(0.02) self.assertFalse(stopping.done()) release.set() await asyncio.wait_for(stopping, 2) tasks = self.engine.store.tasks() - self.assertEqual({t["status"] for t in tasks}, {"succeeded", "pending"}) - worker.im.assert_not_called() + self.assertEqual({t["status"] for t in tasks}, {"succeeded"}) + self.assertEqual(worker.follow.await_count, 1) + self.assertEqual(worker.im.await_count, 1) if __name__ == "__main__": diff --git a/src/test_activity.py b/src/test_activity.py index d0ebb28..df65b98 100644 --- a/src/test_activity.py +++ b/src/test_activity.py @@ -127,17 +127,14 @@ def test_business_stages_skips_dependencies_and_rollback(): assert "1234567890123456789" not in all_logs(audit) store.record_push(main, ["2"]) store.ingest(main, notification("2")) - task = store.claim(worker) - assert task is not None - store.finish(task["id"], "unknown", {"code": "UNCONFIRMED"}, 45) - assert store.claim(worker) is None - before = audit.sequence - assert store.claim(worker) is None - assert audit.sequence == before, "相同依赖等待不应反复刷日志" - store.resolve(task["id"], True) - task = store.claim(worker) - assert task is not None and task["action"] == "dm" - store.finish(task["id"], "succeeded", {"status_code": 0}, 80) + batch = store.claim_batch(worker) + assert {task["action"] for task in batch} == {"follow", "dm"} + follow = next(task for task in batch if task["action"] == "follow") + dm = next(task for task in batch if task["action"] == "dm") + store.finish(follow["id"], "unknown", {"code": "UNCONFIRMED"}, 45) + store.finish(dm["id"], "succeeded", {"status_code": 0}, 80) + assert store.claim_batch(worker) == [] + store.resolve(follow["id"], True) text = all_logs(audit) for stage in ( "通知接收", @@ -150,7 +147,8 @@ def test_business_stages_skips_dependencies_and_rollback(): "任务入队", "任务开始", "任务结果", - "依赖等待", + "冷却开始", + "并行批次", "人工核对", ): assert f"[{stage}]" in text