From ec42517c83a88368728870f0329cc14c76c86b36 Mon Sep 17 00:00:00 2001 From: Rogee Date: Sun, 20 Sep 2026 15:45:05 +0800 Subject: [PATCH] Fix replay of failed history tasks --- ...-09-20-15-31-history-replay-failed-task.md | 33 +++++ src/account_engine.py | 13 ++ src/account_store.py | 127 +++++++++++++++++- src/accounts_app.py | 24 ++-- src/test_accounts.py | 19 +++ 5 files changed, 204 insertions(+), 12 deletions(-) create mode 100644 docs/2026-09-20-15-31-history-replay-failed-task.md diff --git a/docs/2026-09-20-15-31-history-replay-failed-task.md b/docs/2026-09-20-15-31-history-replay-failed-task.md new file mode 100644 index 0000000..ed91b79 --- /dev/null +++ b/docs/2026-09-20-15-31-history-replay-failed-task.md @@ -0,0 +1,33 @@ +# 历史失败任务重新下发 + +- 时间:2026-09-20 15:31 + +## 定位 + +历史事件已经存在于 `events` 表且状态为 `dispatched`。原右键操作再次调用普通 `history_enqueue`,数据库 `INSERT OR IGNORE` 命中同一 `(source,nid)`,不会创建新事件;原失败任务也受任务唯一键和领取冷却保护影响,因此界面没有新的 pending 任务。 + +远端账本查询确认:小号历史任务仍只有原失败记录,没有新增任务记录。将 UID 冷却调整为 1 分钟不能解决事件去重问题。 + +## 修复 + +- 增加专用 `history_replay` 路径,不再重复写入同一历史事件。 +- 只复制原事件中状态为 `failed` 的动作;`succeeded`、`pending`、`running`、`unknown` 不重发。 +- 每个失败动作建立独立 replay event,保留 `replay_of` 指向原事件,原失败记录不改写。 +- 重新下发前要求 UI 明确确认清除目标组内冷却;不支持隐式清除。 +- 新任务使用当前大号规则快照,因此当前 UID 冷却设置会生效。 +- 原小号已删除、未登录、手动停止或跨组时不创建任务。 +- UI 菜单改为“清除冷却并重新执行失败任务”,执行后显示失败数、清除冷却目标数、新建任务数和跳过原因。 + +## 验证 + +- 新增 SQLite 回归测试:确认未显式确认清除时拒绝;确认原失败记录保留;确认 replay 任务可创建;确认重复点击在 pending 期间不重复创建。 +- `test_accounts.py`、`test_plan01.py`、`test_douyin_im.py`、通知及 activity 检查通过。 +- Python 编译、`git diff --check`、LSP error 检查通过。 +- 未执行真实关注、私信或已读操作。 + +## Windows 部署 + +- 已上传源码并在 Windows x64 原生执行 `build_windows.py`。 +- 首次安装因受管 `fingerprint-browser` 进程占用返回 `5`;仅停止该应用路径下的受管浏览器和旧助手进程后重试,安装返回 `0`。 +- 构建目录与安装目录 `DouyinAccounts.exe` SHA-256 一致:`E814C4F679A19ACD705AB5F72F10FC4EB14C1F9FF9393AA26F273A1FD3F5CF4E`。 +- 安装后历史数据仍保留;通过 Windows UI 启动检查标题及界面,未触发真实业务动作。 diff --git a/src/account_engine.py b/src/account_engine.py index 967329b..0ab55d9 100644 --- a/src/account_engine.py +++ b/src/account_engine.py @@ -750,6 +750,17 @@ class Engine: ) return {"selected": len(ids), "events": added, "tasks": tasks} + def replay_history(self, ident, nid): + if not isinstance(nid, str) or not nid.isascii() or not nid.isdecimal(): + raise ValueError("历史事件标识无效") + if self.store.cached_notice(ident, nid) is None: + raise ValueError("历史事件已不在当前本地缓存中") + result = self.store.replay_failed_history( + ident, nid, clear_cooldown=True + ) + result["nid"] = nid + return result + async def fetch_works(self, ident, force=False, session=None, full=False): account = self.store.account(ident) if account["role"] != "main" or account["uid"] is None: @@ -981,6 +992,8 @@ class Engine: return self.enqueue_history( data["id"], data.get("ids"), data.get("baseline") ) + elif name == "history_replay": + return self.replay_history(data["id"], data["nid"]) elif name in ("works_open", "works_refresh", "works_full"): return await self._run_readonly( data["id"], diff --git a/src/account_store.py b/src/account_store.py index 860be41..be18deb 100644 --- a/src/account_store.py +++ b/src/account_store.py @@ -205,6 +205,7 @@ class Store: 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')), + replay_of TEXT NOT NULL DEFAULT '', UNIQUE(source,nid)); CREATE TABLE IF NOT EXISTS tasks ( id INTEGER PRIMARY KEY, event INTEGER NOT NULL REFERENCES events(id), @@ -294,6 +295,10 @@ class Store: self.db.execute( "ALTER TABLE events ADD COLUMN origin TEXT NOT NULL DEFAULT 'live' CHECK(origin IN ('live','history'))" ) + if "replay_of" not in event_columns: + self.db.execute( + "ALTER TABLE events ADD COLUMN replay_of TEXT NOT NULL DEFAULT ''" + ) if "priority" not in task_columns: self.db.execute( "ALTER TABLE tasks ADD COLUMN priority INTEGER NOT NULL DEFAULT 100" @@ -314,7 +319,7 @@ class Store: "UPDATE accounts SET rule=? WHERE id=?", (json.dumps(rule), row["id"]), ) - self.set_setting("schema_version", 3) + self.set_setting("schema_version", 4) 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') @@ -1002,7 +1007,7 @@ class Store: created /= 1000 with self.transaction(): cur = self.db.execute( - "INSERT OR IGNORE INTO events(source,nid,created,targets,rule,state,origin) VALUES (?,?,?,?,?,?,?)", + "INSERT OR IGNORE INTO events(source,nid,created,targets,rule,state,origin,replay_of) VALUES (?,?,?,?,?,?,?,?)", ( source, nid, @@ -1011,6 +1016,7 @@ class Store: json.dumps(rule), "waiting" if enabled else "ignored", origin, + "", ), ) event = self.db.execute( @@ -1075,6 +1081,123 @@ class Store: self._dispatch(source) return bool(cur.rowcount) + def replay_failed_history(self, source, nid, clear_cooldown=False): + if not isinstance(clear_cooldown, bool) or not clear_cooldown: + raise ValueError("重新下发失败任务前必须显式确认清除冷却") + account = self.account(source) + rule = validate_rule(decode(account["rule"])) + if not rule["enabled"]: + raise ValueError("大号规则已关闭,不能重新下发失败任务") + with self.transaction(): + event = self.db.execute( + "SELECT * FROM events WHERE source=? AND nid=?", + (source, nid), + ).fetchone() + if event is None: + raise ValueError("历史事件不存在或已不在本地账本中") + failed = self.db.execute( + "SELECT * FROM tasks WHERE event=? AND status='failed' ORDER BY id", + (event["id"],), + ).fetchall() + if not failed: + raise ValueError("该历史事件没有可重新下发的明确失败任务") + candidates = [] + skipped = [] + for task in failed: + action = task["action"] + if not rule[action]: + skipped.append(f"{action}:当前规则已关闭") + continue + worker = self.account_any(task["worker"]) + if ( + worker["role"] != "worker" + or worker["owner"] != source + or worker["uid"] is None + or worker["task_stopped"] + or worker["delete_state"] + ): + skipped.append(f"{action}:原小号当前不可用") + continue + active = self.db.execute( + "SELECT status FROM tasks WHERE source=? AND target=? AND action=? " + "AND status IN ('pending','running','unknown') ORDER BY id DESC LIMIT 1", + (source, task["target"], action), + ).fetchone() + if active: + skipped.append(f"{action}:已有{active['status']}任务") + continue + candidates.append(task) + if not candidates: + return { + "failed": len(failed), + "cleared": 0, + "tasks": 0, + "skipped": skipped, + } + targets = {task["target"] for task in candidates} + cleared = 0 + for target in targets: + cleared += self.db.execute( + "DELETE FROM cooldowns WHERE source=? AND target=?", + (source, target), + ).rowcount + self.log( + "冷却清除", + "用户显式确认重新下发;已清除失败动作目标的组内冷却", + source=source, + target=target, + ) + replayed = [] + now = time.time() + for task in candidates: + replay_nid = f"{event['nid']}:replay:{task['id']}:{time.time_ns()}" + event_cur = self.db.execute( + "INSERT INTO events(source,nid,created,targets,rule,state,origin,replay_of) VALUES (?,?,?,?,?,?,?,?)", + ( + source, + replay_nid, + event["created"], + json.dumps([task["target"]]), + json.dumps(rule), + "dispatched", + "history", + event["nid"], + ), + ) + task_cur = self.db.execute( + """INSERT INTO tasks(event,source,worker,target,action,params,dependency, + priority,status,result,updated,created_at) + VALUES (?,?,?,?,?,?,NULL,?,'pending','',?,?)""", + ( + event_cur.lastrowid, + source, + task["worker"], + task["target"], + task["action"], + json.dumps(rule), + 0, + now, + now, + ), + ) + replayed.append(task_cur.lastrowid) + self.log( + "任务重新下发", + "仅重下发原历史事件中明确失败的动作;原失败记录保留", + source=source, + worker=task["worker"], + event=event_cur.lastrowid, + task=task_cur.lastrowid, + target=task["target"], + action=task["action"], + ) + return { + "failed": len(failed), + "cleared": cleared, + "tasks": len(replayed), + "skipped": skipped, + } + def dispatch(self, source): with self.transaction(): self._dispatch(source) diff --git a/src/accounts_app.py b/src/accounts_app.py index 2381b80..344cea4 100644 --- a/src/accounts_app.py +++ b/src/accounts_app.py @@ -1397,6 +1397,13 @@ class Window(QMainWindow): "历史事件已确认", f"选择 {payload['selected']} 条;新增事件 {payload['events']} 条;创建任务 {payload['tasks']} 条。\n重复、规则不匹配、作品范围外、冷却或自操作事件不会创建任务。新消息优先。", ) + elif name == "history_replay": + skipped = ";".join(payload.get("skipped") or []) or "无" + QMessageBox.information( + self, + "失败任务重新下发", + f"原明确失败 {payload.get('failed', 0)} 条;清除冷却 {payload.get('cleared', 0)} 个目标;新建任务 {payload.get('tasks', 0)} 条。\n未下发:{skipped}", + ) elif name == "work_filter": text = ( "全部作品" @@ -1490,7 +1497,7 @@ class Window(QMainWindow): return self.history_table.selectRow(row) menu = QMenu(self) - action = menu.addAction("重新执行任务") + action = menu.addAction("清除冷却并重新执行失败任务") action.triggered.connect( lambda _checked=False, nid=nid: self.reexecute_history_event(nid) ) @@ -1511,26 +1518,23 @@ class Window(QMainWindow): QMessageBox.information( self, "重新执行任务预览", - "将重新提交这一条历史事件,按当前规则校验并交给组内小号处理。\n\n" - "可能产生真实关注和私信;已有去重、冷却、占用保护仍然生效,结果不确定不会重发。", + "只重新提交这一条历史事件中明确失败的动作,按当前规则交给原小号处理。\n\n" + "本操作会清除失败动作目标的组内冷却;pending、running、unknown 和已成功动作不会重发。\n" + "可能产生真实关注和私信。", ) if ( QMessageBox.question( self, "确认重新执行任务", - "确认重新提交这条历史事件?", + "确认清除冷却并重新提交失败动作?", QMessageBox.StandardButton.Yes | QMessageBox.StandardButton.No, QMessageBox.StandardButton.No, ) == QMessageBox.StandardButton.Yes ): self.send( - "history_enqueue", - { - "id": payload["account"], - "ids": [nid], - "baseline": payload.get("baseline"), - }, + "history_replay", + {"id": payload["account"], "nid": nid}, ) def works_dialog(self, payload): diff --git a/src/test_accounts.py b/src/test_accounts.py index f1be543..8292c1d 100644 --- a/src/test_accounts.py +++ b/src/test_accounts.py @@ -430,6 +430,25 @@ class StoreTest(unittest.TestCase): {t["status"] for t in self.store.tasks()}, {"failed", "unknown"} ) + def test_failed_history_replay_requires_clear_and_keeps_original_result(self): + self.store.set_rule(self.main, rule(follow=False, dm=True, cooldown=60)) + self.store.ingest(self.main, notice("101", "901", "999"), origin="history") + task = self.store.claim(self.worker) + assert task is not None + self.store.finish(task["id"], "failed", {"code": "CHECK_MSG_NOT_PASS"}) + with self.assertRaises(ValueError): + self.store.replay_failed_history(self.main, "901") + replay = self.store.replay_failed_history(self.main, "901", True) + self.assertEqual(replay["tasks"], 1) + self.assertEqual(replay["cleared"], 1) + tasks = sorted(self.store.tasks(), key=lambda item: item["id"]) + self.assertEqual([item["status"] for item in tasks], ["failed", "pending"]) + replay_event = self.store.db.execute( + "SELECT replay_of FROM events WHERE id=?", (tasks[-1]["event"],) + ).fetchone() + self.assertEqual(replay_event["replay_of"], "901") + self.assertEqual(self.store.replay_failed_history(self.main, "901", True)["tasks"], 0) + def test_disabled_and_cross_group_identity(self): with self.assertRaises(ValueError): self.store.ingest(self.main, notice("102"))