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 拉必然冲突——这份是超集,合一次两份都进。)
This commit is contained in:
+235
@@ -0,0 +1,235 @@
|
||||
"""「已做过」账本:跨设备幂等(评论不重复)+ 进度可见。
|
||||
|
||||
**它解决什么**(用户场景):一台手机登录多个账号、多台手机跑同一个任务,
|
||||
每天跑一次但不知道什么时候跑完,于是反复重跑 → 同一个号被做两次、有的号还没做。
|
||||
账本给出两个东西:
|
||||
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
|
||||
@@ -403,6 +403,39 @@ class TaskStepLog(db.Model):
|
||||
)
|
||||
|
||||
|
||||
class DoneMark(db.Model):
|
||||
"""「已做过」账本:跨设备幂等的标记(「任务 → 去重记录」页)。
|
||||
|
||||
要解决的问题(用户场景):一台手机登录多个账号、多台手机跑同一个任务,
|
||||
任务被反复重跑(因为不知道什么时候跑完)→ 同一个号被做两次、有的号还没做。
|
||||
|
||||
**判据只有一条:`scope_key` 的唯一索引。**
|
||||
多台设备可能同时判断"没做过","先查后插"会两台都插进去;
|
||||
唯一索引 + `INSERT ... ON DUPLICATE KEY`/`INSERT OR IGNORE` 的**受影响行数**
|
||||
才是原子的(见 `core/dedup.py` 的 `mark()`)。
|
||||
|
||||
写入方是任务步骤 `_exec_mark_done`(**成功之后才记账**):动作失败就不记账,
|
||||
下次重跑还会重试该设备——这是"失败不丢"的关键。
|
||||
|
||||
`kind` 是有效期策略(`day`/`hours`/`all`),清理时**只删 day/hours**:
|
||||
`all` 代表"只做一次",删掉就等于去重失效。
|
||||
"""
|
||||
__tablename__ = "done_mark"
|
||||
id = db.Column(db.Integer, primary_key=True, autoincrement=True)
|
||||
scope_key = db.Column(db.String(300), unique=True) # 幂等的全部依据(唯一索引)
|
||||
kind = db.Column(db.String(12), default="day") # day / hours / all
|
||||
job_id = db.Column(db.String(32), default="", index=True)
|
||||
job_name = db.Column(db.String(120), default="")
|
||||
serial = db.Column(db.String(120), default="")
|
||||
device_name = db.Column(db.String(80), default="")
|
||||
identity = db.Column(db.String(200), default="") # 身份值(如抖音号)
|
||||
created_at = db.Column(db.String(20), default="", index=True)
|
||||
|
||||
__table_args__ = (
|
||||
db.Index("ix_done_mark_job_ts", "job_id", "created_at"),
|
||||
)
|
||||
|
||||
|
||||
class AgentConversation(db.Model):
|
||||
"""AI 控制台会话:整个消息序列以 JSON 存在一行里(单会话几十 KB,够用)。"""
|
||||
__tablename__ = "agent_conversation"
|
||||
@@ -427,6 +460,7 @@ SCHEMA_MIGRATIONS = [
|
||||
(4, "自动发现:pending_device 待连接池表(扫描发现的设备,用户确认后才入正式池)", None),
|
||||
(5, "设备池:device 表新增 fingerprint 列(设备指纹 ro.serialno,换 IP 后认领回原记录)", None),
|
||||
(6, "自动发现:pending_device 表新增 fingerprint 列(扫描时读取,用于提示是已有设备换了 IP)", None),
|
||||
(7, "去重账本:done_mark 表(跨设备幂等的「已做过」标记,唯一索引 scope_key)", None),
|
||||
]
|
||||
|
||||
# 当前 schema 版本(备份/恢复用它判断新旧,也写进 app_meta.schema_version)
|
||||
|
||||
@@ -57,6 +57,7 @@ TABLE_LABELS = {
|
||||
"agent_experience": "经验库", "experience_audit": "经验巡检",
|
||||
"agent_action": "动作库", "device_install_log": "设备端安装记录",
|
||||
"task_step_log": "任务步骤明细",
|
||||
"done_mark": "去重记录(已做过)",
|
||||
}
|
||||
_STAGE_TTL = 1800 # 导入暂存有效期(秒)
|
||||
|
||||
|
||||
+17
-4
@@ -148,8 +148,10 @@ def validate_steps(steps, depth=1, path="steps", errors=None, warnings=None,
|
||||
errors.append(f"{here}.params.probability: 必须是 0~100 的数字"
|
||||
f"(当前 {prob!r})")
|
||||
|
||||
# 选择器类步骤
|
||||
if stype in NEED_SELECTOR:
|
||||
# 选择器类步骤(例外:if_el 用「去重」条件时不需要 selector_value,
|
||||
# 它的身份元素在 ident_value 里,按自己的规则校验)
|
||||
if stype in NEED_SELECTOR and not (stype == "if_el"
|
||||
and params.get("selector_type") == "dedup"):
|
||||
sel = params.get("selector_value")
|
||||
if not isinstance(sel, str) or not sel.strip():
|
||||
errors.append(f"{here}.params.selector_value: 不能为空 —— "
|
||||
@@ -234,6 +236,17 @@ def validate_steps(steps, depth=1, path="steps", errors=None, warnings=None,
|
||||
elif (params.get("selector_type") or "xpath") == "screen":
|
||||
warnings.append(f"{here}: 屏幕状态没有文本可比,"
|
||||
"cmp_op/cmp_value 会被忽略")
|
||||
# 去重条件(见 tasks/generic/task.py 的 selector_type="dedup"):
|
||||
# 不用 selector_value,身份元素填在 ident_value
|
||||
if (params.get("selector_type") or "") == "dedup":
|
||||
iv = (params.get("ident_value") or "").strip()
|
||||
if not iv:
|
||||
warnings.append(f"{here}: 去重条件没填身份元素,将按「设备」当身份"
|
||||
"(一号一机时没问题;一台机器多个号时要填账号那个元素)")
|
||||
elif (params.get("ident_type") or "text") == "xpath" \
|
||||
and not (iv.startswith("//") or iv.startswith("(//")):
|
||||
errors.append(f"{here}.params.ident_value: 身份元素的类型是 xpath,"
|
||||
f"取值应以 // 或 (// 开头(当前 {iv[:40]!r})")
|
||||
|
||||
if stype == "swipe":
|
||||
_check_direction(params, here, errors)
|
||||
@@ -294,9 +307,9 @@ def _check_selector(params, here, stype, errors, warnings):
|
||||
stype_ok = IF_SELECTOR_TYPES if stype == "if_el" else SELECTOR_TYPES
|
||||
sel_type = params.get("selector_type") or "xpath"
|
||||
value = (params.get("selector_value") or "").strip()
|
||||
if sel_type not in stype_ok:
|
||||
if sel_type not in stype_ok and not (sel_type == "dedup" and stype == "if_el"):
|
||||
errors.append(f"{here}.params.selector_type: 不支持 {sel_type!r}"
|
||||
f"(可用:{'、'.join(stype_ok)})")
|
||||
f"(可用:{'、'.join(stype_ok + ('dedup',))})")
|
||||
return
|
||||
if len(value) > MAX_SELECTOR_LEN:
|
||||
errors.append(f"{here}.params.selector_value: 太长(>{MAX_SELECTOR_LEN} 字符)")
|
||||
|
||||
Reference in New Issue
Block a user