用户场景(他原话):一台手机登录 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 拉必然冲突——这份是超集,合一次两份都进。)
236 lines
10 KiB
Python
236 lines
10 KiB
Python
"""「已做过」账本:跨设备幂等(评论不重复)+ 进度可见。
|
||
|
||
**它解决什么**(用户场景):一台手机登录多个账号、多台手机跑同一个任务,
|
||
每天跑一次但不知道什么时候跑完,于是反复重跑 → 同一个号被做两次、有的号还没做。
|
||
账本给出两个东西:
|
||
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
|