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