Fix replay of failed history tasks
This commit is contained in:
@@ -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 启动检查标题及界面,未触发真实业务动作。
|
||||
@@ -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"],
|
||||
|
||||
+125
-2
@@ -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)
|
||||
|
||||
+14
-10
@@ -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):
|
||||
|
||||
@@ -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"))
|
||||
|
||||
Reference in New Issue
Block a user