Files
auto_control/core/dedup.py
T
butubb 876224f876 feat(去重): 跨设备「已做过」账本 —— 同一个号不会做两次 + 「谁做过了」看得见
用户场景(他原话):一台手机登录 5 个抖音号、一共 5 台手机,每个任务只让其中一个
目标号评论;每天跑一次但不知道什么时候跑完,于是"一直重复跑" → 结果
"一个手机还没评论到,一个手机都评论两次了"。

**根因不是"单设备重复",是跨设备没有共享的判断 + 进度不可见。** 所以做两件事:
① 幂等;② 把"谁做过了、还差谁"摆到台面上(不然只能靠重跑确认,而重跑又在制造重复)。

- `core/models.py`:新表 `done_mark`(迁移账本补 v7)。**判据只有 `scope_key` 的
  唯一索引**——多台设备会同时判断"没做过","先查后插"有竞态(两台都插),
  唯一索引 + `INSERT ... ON DUPLICATE KEY`/`INSERT OR IGNORE` 的**受影响行数**才原子。
- `core/dedup.py`(新):`build_key`(`任务|身份|时间桶`)/ `check` / `mark` /
  `list_marks`(带"今天做了几台/几个号"统计)/ `delete_mark` / `clear_job` / `purge_old`。
  自建 app context(照 device_pool 的 `_ctx()`),任务线程/Web/清理都不用关心。
- 任务侧两个部件(**检查在前、记账在后**):
  · `if_el` 新增条件类型 `selector_type="dedup"`:命中=这个身份做过了 → 走 then 分支。
    身份元素在 `ident_type`/`ident_value`(留空 = 用设备 serial,一号一机场景)。
  · 新步骤 `mark_done`「记为已做」(22 种步骤):放动作**成功之后**。
  拆两步的用意:动作失败就不记账,下次重跑还会重试该设备 —— 失败不丢。
- 有效期(`dedup_reset` = day/all/hours)放**任务级**:检查与记账两处各填一份的话,
  填不一致就算出两个 key、去重会**静默失效**,所以强制只配一处(编辑器顶部下拉)。
- 三条防误伤规则(都有测试兜着):
  · 身份读不到 / 身份值过长 → **不去重、当没做过照常执行**。绝不能把"读不到"
    当成空身份——那会让所有设备共用一个 key、第一台记账后其余全被误判成"做过"。
  · `kind='all'`(只做一次)的记录**永不清理**(清了等于语义失效);清理只删 day/hours。
  · 去重的两个易错点在保存时直接告警:身份元素两边不一致、有检查没记账/有记账没检查。
- 「任务 → 去重记录」新子分栏(`static/admin/dedup.js`):统计行 + 明细表 +
  删单条(那个号重跑)/ 清空任务(整批重跑)。接口 3 个(GET/delete/clear,PERM_TASKS)。
- 每日 04:23 清理(挂现有 APScheduler),`TABLE_LABELS` 补中文名(备份覆盖自动派生)。
- AI 建任务草稿校验同步:`dedup` 走自己的规则(要 ident_value、xpath 前缀校验),
  没填身份元素只警告不拦(用设备当身份是合法用法);普通条件空选择器仍然拦。
- 文档:TASK_DEV §4.6(去重专章 + App 内检测的兜底配方与它的三个局限)、
  DATA_MODEL §2.9、API 三个接口、ARCHITECTURE(分层/装配/子分栏/JS 分工/清理)、
  DEPLOY §5.2(15 张表)、步骤数 21→22 全库同步。

自测:单元 + 集成 33 项(**含 8 线程抢同一个身份、恰好一个成功**的原子性断言,
以及"all 记录不被清理""身份读不到不去重""清了能重跑")、
**真机端到端**(cs1 上"检查→动作→记账"跑两遍:第二遍被拦、换 serial 的"另一台设备"
同样被拦、删记录后能重跑)、草稿校验 5 项、GET 冒烟 56 路由 0 个 500。

(注:本分支基于 feat/if-el-multi-value,因为它俩都要改 task.py 的 STEP_TYPES 与
editor.js 的 STEP_LIB 同一区域,分开从 dev 拉必然冲突——这份是超集,合一次两份都进。)
2026-09-24 10:55:29 +08:00

236 lines
10 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""「已做过」账本:跨设备幂等(评论不重复)+ 进度可见。
**它解决什么**(用户场景):一台手机登录多个账号、多台手机跑同一个任务,
每天跑一次但不知道什么时候跑完,于是反复重跑 → 同一个号被做两次、有的号还没做。
账本给出两个东西:
1. `check()` —— "这个身份在这个任务里做过了吗"(跨设备,多台手机共享同一份判断)
2. `list_marks()` —— "谁做过了、还差谁"(界面上看得见,就不用靠"一直重跑"来确认)
**为什么必须靠唯一索引**:多台设备可能同时判断"没做过"。
"先 SELECT 再 INSERT"有竞态——两台都会插进去。唯一索引 + `INSERT ... ON DUPLICATE KEY`
(MySQL)/ `INSERT OR IGNORE`(SQLite)的**受影响行数**才是原子的判据:
`1` = 你抢到了首次,`0` = 别人已经做过。
**身份值读不到时绝不去重**(`identity is None`):宁可漏拦一次,
也不能让"读不到"退化成"空身份"——那会让所有设备共用一个 key,
第一台记账后其余全部被误跳过。同理,身份值过长时也不截断(截断会让两个身份撞车)。
有效期策略(`kind`)由任务级参数 `dedup_reset` 决定,**只在一处配**:
步骤里各填一份的话,两处填不一致就会算出不用的 key、去重静默失效。
DB 访问由本模块自建 app context(照 `core/device_pool.py` 的 `_ctx()`),
调用方(任务线程 / Web 线程 / 每日清理)都不用关心。
"""
import time
from datetime import datetime, timedelta
from core.logger import get_logger
from core.models import db
_log = get_logger("core.dedup")
# 有效期策略
KIND_DAY = "day" # 每天一次(默认;"每天跑一次"的场景)
KIND_HOURS = "hours" # 每 N 小时一次
KIND_ALL = "all" # 只做一次(永久)
RESET_KINDS = (KIND_DAY, KIND_ALL, KIND_HOURS)
DEFAULT_HOURS = 6
# 账本保留期:只影响 day/hours 桶(它们过期就没意义了);
# **kind='all' 的记录永不清理**——清了就等于"只做一次"失效
KEEP_DAYS = 180
# scope_key 上限(= DoneMark.scope_key 列宽)。超了不去重、只告警:
# 截断会让两个不同身份撞成同一个 key,进而误跳过
MAX_KEY_LEN = 290
_app = None
def init_app(app):
"""web_server 启动时调用:绑 app(后台线程访问 db 要自推 context)。"""
global _app
_app = app
def _ctx():
if _app is None:
raise RuntimeError("dedup 未关联 Flask app(web_server 启动时调用 init_app)")
return _app.app_context()
# ================== 唯一键 ==================
def build_key(job_id, identity, kind=KIND_DAY, hours=DEFAULT_HOURS, now=None):
"""算唯一键:`任务 | 身份 | 时间桶`。
桶只影响"多久之后算新的一轮":
· day → `|d:2026-09-24`(本地日期,跟着自然日走)
· hours → `|h6:4939219`(epoch // (N*3600),改 N 就换了桶)
· all → `|all`(无桶,永远命中同一个 key)
"""
now = time.time() if now is None else now
base = f"{(job_id or '').strip()}|{(identity or '').strip()}"
if kind == KIND_ALL:
return base + "|all"
if kind == KIND_HOURS:
h = max(1, int(hours or DEFAULT_HOURS))
return f"{base}|h{h}:{int(now // (h * 3600))}"
return base + "|d:" + time.strftime("%Y-%m-%d", time.localtime(now))
def _key_or_none(job_id, identity, kind, hours):
"""算 key 并做长度体检;超长返回 None(调用方按"不去重"处理,不截断)。"""
if not (identity or "").strip():
return None
key = build_key(job_id, identity, kind, hours)
if len(key) > MAX_KEY_LEN:
_log.warning(f"去重键过长({len(key)} 字符)→ 本次不去重,避免截断后两个身份撞车: {key[:80]}…")
return None
return key
# ================== 查 / 记 ==================
def check(job_id, identity, kind=KIND_DAY, hours=DEFAULT_HOURS):
"""这个身份在这个任务里做过了吗?(只读)
返回 True=做过 / False=没做过。identity 为空时返回 False(不去重,照常执行)。
"""
key = _key_or_none(job_id, identity, kind, hours)
if key is None:
return False
try:
with _ctx():
row = db.session.execute(
db.text("SELECT 1 FROM done_mark WHERE scope_key = :k LIMIT 1"),
{"k": key}).fetchone()
return row is not None
except Exception as e:
# 账本读失败不能拦住任务:当"没做过"放行(宁可重复,不可卡死)
_log.warning(f"去重检查失败(按未做过放行): {e}")
return False
def mark(job_id, identity, kind=KIND_DAY, hours=DEFAULT_HOURS,
serial="", device_name="", job_name=""):
"""原子记账。返回 True=本次是首次(你抢到了)/ False=别人已经记过(或没记成)。
⚠ **必须靠唯一索引**:不能先 check 再 insert(多台设备会同时通过)。
"""
key = _key_or_none(job_id, identity, kind, hours)
if key is None:
return False
args = {"k": key, "kind": kind, "jid": job_id or "", "jname": job_name or "",
"serial": serial or "", "dname": device_name or "",
"ident": (identity or "").strip()[:200],
"ts": time.strftime("%Y-%m-%d %H:%M:%S")}
try:
with _ctx():
is_mysql = db.engine.dialect.name == "mysql"
cols = ("scope_key, kind, job_id, job_name, serial, device_name, "
"identity, created_at")
vals = (":k, :kind, :jid, :jname, :serial, :dname, :ident, :ts")
if is_mysql:
# 冲突时做一次"无变化"的更新:受影响行数 1=插入、0=已存在(MySQL 语义)
sql = (f"INSERT INTO done_mark ({cols}) VALUES ({vals}) "
f"ON DUPLICATE KEY UPDATE scope_key = scope_key")
else:
sql = f"INSERT OR IGNORE INTO done_mark ({cols}) VALUES ({vals})"
r = db.session.execute(db.text(sql), args)
db.session.commit()
first = (r.rowcount or 0) > 0
if first:
_log.info(f"[{serial}] 记账:任务『{job_name or job_id}』身份『{args['ident']}』"
f"({kind})")
return first
except Exception as e:
_log.warning(f"去重记账失败(本次不记,下次重跑会重试): {e}")
return False
# ================== 界面:看 / 清 ==================
def list_marks(job_id="", limit=200):
"""去重记录(界面用):返回 (rows, stats)。
stats:`total` 本任务全部 · `today_devices` 今天已做的设备数 ·
`today_identities` 今天已做的身份值数(== "今天做成了几个号")。
"""
limit = max(1, min(int(limit or 200), 2000))
today0 = time.strftime("%Y-%m-%d 00:00:00")
where, args = "", {}
if job_id:
where = "WHERE job_id = :jid"
args["jid"] = job_id
with _ctx():
rows = db.session.execute(db.text(
f"SELECT id, scope_key, kind, job_id, job_name, serial, device_name, "
f"identity, created_at FROM done_mark {where} "
f"ORDER BY id DESC LIMIT :lim"), {**args, "lim": limit}).fetchall()
twhere = where + (" AND " if where else "WHERE ") + "created_at >= :t0"
stat = db.session.execute(db.text(
f"SELECT COUNT(*) AS total, "
f"COUNT(DISTINCT serial) AS devs, "
f"COUNT(DISTINCT identity) AS ids "
f"FROM done_mark {where}"), args).fetchone()
today = db.session.execute(db.text(
f"SELECT COUNT(DISTINCT serial) AS devs, "
f"COUNT(DISTINCT identity) AS ids "
f"FROM done_mark {twhere}"), {**args, "t0": today0}).fetchone()
out = [{"id": r[0], "scope_key": r[1], "kind": r[2], "job_id": r[3], "job_name": r[4],
"serial": r[5], "device_name": r[6], "identity": r[7], "created_at": r[8]}
for r in rows]
stats = {"total": stat[0] if stat else 0,
"devices": stat[1] if stat else 0,
"identities": stat[2] if stat else 0,
"today_devices": today[0] if today else 0,
"today_identities": today[1] if today else 0}
return out, stats
def delete_mark(mark_id):
"""删一条记录(让某个号/某台设备能重跑)。返回是否删掉了。"""
try:
with _ctx():
r = db.session.execute(db.text("DELETE FROM done_mark WHERE id = :i"),
{"i": int(mark_id)})
db.session.commit()
return (r.rowcount or 0) > 0
except Exception as e:
_log.warning(f"删除去重记录失败: {e}")
return False
def clear_job(job_id):
"""清空某个任务的全部去重记录(整批重跑)。返回删掉的行数。"""
try:
with _ctx():
r = db.session.execute(db.text("DELETE FROM done_mark WHERE job_id = :j"),
{"j": job_id or ""})
db.session.commit()
n = r.rowcount or 0
_log.info(f"清空去重记录:任务 {job_id} 共 {n} 条")
return n
except Exception as e:
_log.warning(f"清空去重记录失败: {e}")
return 0
def purge_old(keep_days=KEEP_DAYS):
"""清理过期的 day/hours 桶记录。**kind='all' 永不清理**(清了等于去重失效)。
量级很小(设备数 × 天数),直接一条 DELETE 就行,不用像步骤明细那样分批。
"""
cutoff = (datetime.now() - timedelta(days=max(1, int(keep_days)))
).strftime("%Y-%m-%d %H:%M:%S")
try:
with _ctx():
r = db.session.execute(db.text(
"DELETE FROM done_mark WHERE kind <> :all AND created_at < :cut"),
{"all": KIND_ALL, "cut": cutoff})
db.session.commit()
n = r.rowcount or 0
if n:
_log.info(f"去重记录清理:删除 {n} 条(保留 {keep_days} 天;"
f"kind={KIND_ALL} 的永不清理)")
return n
except Exception as e:
_log.warning(f"去重记录清理失败: {e}")
return 0