Files
auto_control/core/device_discovery.py
T
butubb 97cec6e211 feat(通知): 系统级 Webhook 通知子系统(事件目录 + 可插拔适配器 + 防刷屏)
平台此前出了问题只能靠人盯页面。现在各组件统一走 `notifier.notify(事件, **字段)`,
推到企业微信 / 自建服务;**所有可通知点都登记进事件目录,默认全关,用户按 webhook 勾选**。

## 架构(`core/notifier.py` + `core/notify_events.py`)
业务线程调 `notify()` → 只做内存操作(读配置快照/匹配订阅/入队)→ 返回;
后台 1 个 dispatcher(聚合 + 每 hook 限流 + 折叠摘要)+ 3 个 sender(真实 HTTP、退避重试)
负责真正发出去。硬约束:**notify 零 DB、零 HTTP、零阻塞、异常不冒泡**——所以任务线程里
可以直接调(不用 app_context、不用 try/except),但**必须放在所有 `with self._lock` 之外**。

- **事件目录 36 条**(任务批次/单设备/设备/Worker/业务/安装/系统/AI),支持 `task.*` 通配订阅;
  语义分工避免重复告警:`worker.*` 是单次尝试级,`task.device.success/failed` 是唯一权威结论。
- **适配器可插拔**:`wecom`(markdown,按 4096 **字节**截断、超限不截半个汉字)+
  `json`(模板占位符,替换值按 JSON 转义,保存前干跑校验);钉钉/飞书留了插槽(前端置灰)。
- **防打爆四层**:聚合窗口(默认 30s,同批次合并成一条并带样本)→ 令牌桶限流(默认 18/分,
  对齐企微硬限 20)→ 被限流的**折叠成摘要不丢弃** → 有界队列背压。取舍:失败通知最多延迟
  一个窗口,换来群不被刷屏。

## 安全与存储
- 配置只落 `app_meta.notify_webhooks` 一个键(**不建表** → 不涉及备份覆盖红线)。
- URL 本身就是凭据(企微 `?key=`)→ 接口回显/发送记录/日志一律 `mask_url()/scrub()`;
  编辑时留空即不修改;secret 永不回显。DATA_MODEL 的明文凭据告警补上了这一条。
- 发送记录:内存环形缓冲 200 条(重启清空)+ 独立 `logs/notify.log`。

## 接入点(每个都放在锁外、不改 return 顺序)
task_manager(批次开始/结束用新增的 `_BatchTracker` 统一在 finally 计数、单设备成功/失败/
离线/重试/停止/抢占/归还/cron 停止)、device_worker 心跳看门狗、generic 任务选择器连续失效、
apk 安装开始/完成、设备上下线(**状态沿检测**,只报新变化)、备份导出/恢复、经验巡检、
用户登录、服务启停。

## 前端
系统 Tab 新增「通知 / Webhook」子分栏:多条 webhook 列表(URL 打码)+ 编辑弹窗(格式/URL/
密钥/事件勾选树带 ★建议/聚合/限流/自定义模板/预览)+ 发送测试 + 发送记录。

## 自测
- 进程内逻辑 10 组断言全绿:聚合合并、限流+折叠、无配置/全局关静默丢弃、未知事件、
  内部异常不外泄、JSON 转义(标题含引号换行仍合法)、URL/异常消息脱敏、配置校验。
- 端到端(假 webhook 接收端)17 项断言全绿:真实事件投递(user.login / task.batch.no_device)、
  企微请求体形状、**HTTP 200 + errcode 93000 判为失败**、500 重试 3 次、记录里 URL 打码。
- 韧性:webhook 指向黑洞地址时登录耗时 100~114ms(基线 107~133ms,**异步隔离生效**);
  配置写成坏 JSON 服务照常启动、通知静默不发、日志有 error(服务端实测后已复原)。
- 页面:系统 → 通知 面板/弹窗/36 个事件复选框/预览全部正常,无 JS 报错。
- 自测数据已清理(webhook、自建任务、写坏又复原的配置键)。

文档:新增 doc/NOTIFY.md(事件表/配置/格式约束/防刷屏/加事件三步骤/排障)并登记进 doc/README;
API.md §2.11;DATA_MODEL 的 app_meta 键表与明文凭据告警;ARCHITECTURE 线程表/分层/扩展点;
根 README 功能索引与日志表。
2026-09-15 13:55:14 +08:00

507 lines
21 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.
"""设备自动发现:扫描网段中开放 adb 5555 的设备 → 待连接池(pending)。
流程:socket 并发探测 5555(标准库,不碰 adb)→ 剔除已在设备池的 serial →
串行 adb 验证(adb_connect_light,内部全局锁)→ 用 `adb devices` 的
state=="device" 过滤(排除 unauthorized/offline,connect 输出不可信)→
写入 pending_device 表。**扫描只验证、不自动连接入池**——用户在前端确认后
才调 device_pool.add_device + adb_connect(见 web/devices_api.py)。
安全红线:绝不 kill-server / 绝不 disconnect(与全项目一致)。
配置存 app_meta KV 表(discovery_enabled/subnets/interval/port),定时线程
每轮重读,开关/周期/网段改动即时生效无需重启。
"""
import ipaddress
import json
import socket
import threading
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
from config import (USB_ADB_HOST,
DISCOVERY_PORT, DISCOVERY_SUBNETS, DISCOVERY_INTERVAL)
from core.adb_helper import _adb, adb_connect_light
from core.logger import get_logger
from core import notifier
from core.models import db, Device, PendingDevice
_log = get_logger("core.disc")
# 扫描配置键(app_meta)
_K_ENABLED = "discovery_enabled"
_K_SUBNETS = "discovery_subnets"
_K_INTERVAL = "discovery_interval"
_K_PORT = "discovery_port"
# 指纹匹配时是否自动认领(默认关:认领会改写分组/任务引用,默认交人工确认)
_K_AUTO_CLAIM = "discovery_auto_claim"
# 单网段主机数上限(防误配 /8 之类),超过截断并告警
_MAX_HOSTS_PER_SUBNET = 1024
# socket 探测并发数与单 IP 超时
_PROBE_WORKERS = 64
_PROBE_TIMEOUT = 0.4
_app = None
_scan_lock = threading.Lock() # 定时/手动扫描互斥
# 上一轮"池内离线"集合:通知只报**状态沿**(这一轮新变成离线的),不每轮刷屏
_last_offline = set()
_stop_event = threading.Event() # shutdown 用
_scanning = False # 状态快照(API/前端)
_last_scan = None # (时间串, 开放数, 可连数, 新增数)
_last_error = ""
def _ctx():
"""后台线程访问 db 时自行推 app context(仿 device_pool._ctx)。"""
return _app.app_context() if _app else None
def _fmt(ts=None):
return (ts or datetime.now()).strftime("%Y-%m-%d %H:%M")
# ================== 配置读写 ==================
def get_settings():
"""读发现配置(app_meta,缺键用 config 默认值补齐)。"""
with _ctx():
from core.db_config import meta_get
def _get(key, default):
v = meta_get(key)
return v if v is not None else default
try:
subnets = json.loads(_get(_K_SUBNETS, "[]")) or DISCOVERY_SUBNETS
except Exception:
subnets = DISCOVERY_SUBNETS
return {
"enabled": _get(_K_ENABLED, "1") == "1", # 默认开启
"subnets": subnets,
"interval": int(_get(_K_INTERVAL, str(DISCOVERY_INTERVAL)) or DISCOVERY_INTERVAL),
"port": int(_get(_K_PORT, str(DISCOVERY_PORT)) or DISCOVERY_PORT),
"auto_claim": _get(_K_AUTO_CLAIM, "0") == "1", # 默认关
}
def save_settings(enabled=None, subnets=None, interval=None, port=None, auto_claim=None):
"""写发现配置(部分字段更新)。subnets 逐项校验 CIDR,非法返回 (False, 原因)。"""
if subnets is not None:
clean = []
for s in subnets:
s = (s or "").strip()
if not s:
continue
try:
ipaddress.ip_network(s, strict=False)
except ValueError:
return False, f"网段格式错误: {s}"
clean.append(s)
subnets = clean
try:
if interval is not None and (int(interval) < 10 or int(interval) > 3600):
return False, "扫描周期需在 10-3600 秒之间"
if port is not None and (int(port) < 1 or int(port) > 65535):
return False, "端口不合法"
except (TypeError, ValueError):
return False, "参数不合法"
with _ctx():
from core.db_config import meta_set
_put = meta_set # 方言中立 upsert(app_meta.key 在 MySQL 里是保留字)
if enabled is not None:
_put(_K_ENABLED, "1" if enabled else "0")
if subnets is not None:
_put(_K_SUBNETS, json.dumps(subnets))
if interval is not None:
_put(_K_INTERVAL, str(int(interval)))
if port is not None:
_put(_K_PORT, str(int(port)))
if auto_claim is not None:
_put(_K_AUTO_CLAIM, "1" if auto_claim else "0")
db.session.commit()
return True, "已保存"
# ================== 网段展开与端口探测 ==================
def _local_ips():
"""本机自身 IP 集合(尽力而为):扫描时排除,避免探测到自己的 5555。
优先 `hostname -I`(Linux/macOS 支持,一行空格分隔多 IP);不支持/失败的
平台退回 socket.getaddrinfo 枚举。关键:外部命令必须用 bytes 收——Windows
上 git-bash 的 coreutils hostname 不支持 -I,会把 GBK 报错写进 stderr,
text=True 在 subprocess 后台读线程里 utf-8 严格解码会直接炸线程(主线程
try/except 接不住异步线程异常)。
"""
import subprocess
ips = set()
try:
r = subprocess.run(["hostname", "-I"], capture_output=True, timeout=3)
if r.returncode == 0:
ips.update(p for p in
r.stdout.decode("utf-8", errors="ignore").split() if p)
except Exception:
pass
if not ips: # 兜底:主机名解析出的接口 IPv4
try:
for info in socket.getaddrinfo(socket.gethostname(), None):
ip = info[4][0]
if ":" not in ip:
ips.add(ip)
except Exception:
pass
return ips
def _expand_subnets(subnets, max_hosts=_MAX_HOSTS_PER_SUBNET):
"""CIDR 列表 → IP 列表。排除 220 自身(USB_ADB_HOST)与本机 IP;
非法网段跳过记日志;单网段超过 max_hosts 截断并告警。"""
self_ips = {USB_ADB_HOST} | _local_ips()
ips = []
for cidr in subnets:
try:
net = ipaddress.ip_network((cidr or "").strip(), strict=False)
except ValueError:
_log.warning(f"发现: 跳过非法网段 {cidr}")
continue
hosts = [str(h) for h in net.hosts()
if str(h) not in self_ips]
if len(hosts) > max_hosts:
_log.warning(f"发现: 网段 {cidr} 主机数 {len(hosts)} 超过上限,截断前 {max_hosts}")
hosts = hosts[:max_hosts]
ips.extend(hosts)
return ips
def _probe_port(ip, port, timeout=_PROBE_TIMEOUT):
"""单 IP 端口探测(标准库 socket,超时/拒绝/不可达一律 False)。"""
try:
with socket.create_connection((ip, port), timeout=timeout):
return True
except Exception:
return False
def _probe_open(ips, port, workers=_PROBE_WORKERS):
"""并发探测,返回开放端口的主机列表。"""
if not ips:
return []
with ThreadPoolExecutor(max_workers=workers) as ex:
results = list(ex.map(lambda ip: (ip, _probe_port(ip, port)), ips))
return [ip for ip, ok in results if ok]
# ================== adb devices 解析 ==================
def _parse_adb_devices(out):
"""解析 `adb devices` 输出 → state=="device" 的 serial 集合。
只看 connect 输出不可靠(未授权设备也返回 connected),必须用 state 过滤。
"""
devices = set()
for line in (out or "").splitlines()[1:]:
parts = line.split()
if len(parts) >= 2 and parts[0] and not parts[0].startswith("*"):
if parts[1] == "device":
devices.add(parts[0])
return devices
# ================== 自动认领(可选) ==================
def _auto_claim(fps):
"""指纹命中池中已有设备 → 自动认领:迁移记录到新地址并同步分组/任务引用。
仅在设置 discovery_auto_claim 打开时由扫描线程调用(默认关)。
返回 [(old_serial, new_serial, name), ...]。
"""
from core import device_pool
done = []
for serial, fp in (fps or {}).items():
if not fp:
continue
try:
info = device_pool.find_by_fingerprint(fp)
if not info or info.get("serial") == serial:
continue # 没匹配到,或本来就是这条(无需迁移)
old, name = device_pool.claim_device(serial, fp)
if not old:
continue
PendingDevice.query.filter_by(serial=serial).delete()
db.session.commit()
done.append((old, serial, name or info.get("name") or ""))
_log.info(f"自动认领: 『{name or info.get('name') or old}』{old} → {serial}")
except Exception as e:
db.session.rollback()
_log.warning(f"自动认领 {serial} 失败: {e}")
return done
# ================== 扫描 ==================
def scan_once(manual=False):
"""执行一轮扫描。返回 (ok, result_dict);后台线程/API 调用。"""
global _scanning, _last_scan, _last_error
if not _scan_lock.acquire(blocking=False):
return False, {"error": "扫描进行中"}
_scanning = True
try:
settings = get_settings()
port = settings["port"]
subnets = settings["subnets"]
if not subnets:
_last_error = "未配置扫描网段"
return True, {"found": 0, "verified": 0, "new": 0, "error": _last_error}
# 1. socket 并发探测(不碰 adb,把几百台缩小到开放端口量级)
open_ips = _probe_open(_expand_subnets(subnets), port)
candidates = {f"{ip}:{port}" for ip in open_ips}
with _ctx():
# 2. 剔除已在正式设备池的 serial
configured = set(device_pool_list_configured())
candidates -= configured
# 3. 串行 adb 验证:已连接的直接跳过(零成本),其余 adb_connect_light
already = _parse_adb_devices(_adb("devices"))
for serial in candidates - already:
adb_connect_light(serial) # 内部全局锁串行;失败无妨,下面按 state 过滤
# 4. state==device 过滤(排除 unauthorized/offline)
verified = _parse_adb_devices(_adb("devices")) & candidates
# 5. 写 pending:新 → 插入;已有 → 更新 last_seen(顺带刷新指纹)
# 指纹用于认出"这台其实是设备池里某台设备换了 IP",见 list_pending 的 match
now = _fmt()
existing = {p.serial for p in PendingDevice.query.all()}
added = 0
added_serials = []
from core import device_pool
fps = {}
for serial in verified:
source = "tailscale" if serial.split(":")[0].startswith("100.") else "lan"
try:
fp = device_pool.read_fingerprint(serial, timeout=4)
except Exception:
fp = ""
fps[serial] = fp
if serial in existing:
PendingDevice.query.filter_by(serial=serial).update(
{"last_seen": now, "fingerprint": fp})
else:
db.session.add(PendingDevice(serial=serial, source=source,
first_seen=now, last_seen=now,
fingerprint=fp))
added += 1
added_serials.append(serial)
db.session.commit()
# 5.5 自动认领(可选,默认关):指纹命中池中已有设备 → 直接把记录迁到新地址。
# 默认关是因为认领会改写分组/任务引用(数据结构变动),交人工点一下更稳妥;
# 打开后零点击完成,见 doc/API.md §6。
claimed = _auto_claim(fps) if settings.get("auto_claim") else []
# 6. 正式池断联设备自动重连:adb connect 会因 WiFi 波动/设备重启/
# adb 服务重启而断开——扫描线程每轮顺带重试(幂等轻量,内部
# 全局锁串行),连上即恢复在线,无需人工干预。pending 池是给
# 「未授权新设备」的,正式池设备断联不进 pending,而是自动重连。
back = _reconnect_offline(configured)
# 7. 离线集合(供"状态沿"通知用:只报**这一轮新变成离线**的,不每轮刷屏)
offline_now = {s for s in configured if not device_pool.is_online(s)}
_log.info(f"发现: 探测开放 {len(open_ips)} 台,可连 {len(verified)} 台,"
f"新增待连接 {added} 台"
+ (f",自动重连恢复 {len(back)} 台 {back}" if back else "")
+ (f",自动认领 {len(claimed)} 台 {[c[0] + '→' + c[1] for c in claimed]}"
if claimed else ""))
# 通知统一放在 `with _ctx():` 之后(别在 DB 事务期间做额外的事)
global _last_offline
try:
newly_offline = sorted(offline_now - _last_offline)
_last_offline = set(offline_now)
if newly_offline:
notifier.notify("device.offline", serials=newly_offline,
count=len(newly_offline),
devices=newly_offline[:10])
if back:
notifier.notify("device.online", serials=list(back), count=len(back))
if added_serials:
notifier.notify("device.discovered", serials=added_serials[:10],
count=len(added_serials))
if claimed:
notifier.notify("device.claimed",
pairs=[f"{c[0]}→{c[1]}" for c in claimed][:10],
count=len(claimed))
except Exception as e:
_log.warning(f"发送设备上下线通知失败(不影响扫描): {e}")
_last_scan = (now, len(open_ips), len(verified), added)
_last_error = ""
return True, {"found": len(open_ips), "verified": len(verified),
"new": added, "claimed": len(claimed)}
except Exception as e:
_log.warning(f"发现扫描异常: {e}")
_last_error = str(e)[:200]
return False, {"error": _last_error}
finally:
_scanning = False
_scan_lock.release()
def device_pool_list_configured():
"""设备池已配置 serial(延迟 import 避免循环依赖)。"""
from core import device_pool
return device_pool.list_configured()
# ================== 正式池断联设备:自动重连 ==================
def _reconnect_offline(configured):
"""对正式池中断联的网络设备逐个 adb 重连,返回恢复的 serial 列表。
只重连网络设备(IP:5555;USB 设备插着就在,无需 connect)。
幂等轻量:内部 adb 全局锁串行,失败静默(下轮扫描再试)。
"""
online_now = _parse_adb_devices(_adb("devices"))
targets = [s for s in (configured or [])
if ":" in s and s not in online_now]
if not targets:
return []
for serial in targets:
adb_connect_light(serial)
online_after = _parse_adb_devices(_adb("devices"))
return [s for s in targets if s in online_after]
def list_pool_offline():
"""正式设备池中断联的设备(serial + 型号 + 备注名),面板展示用。
断联设备仍是正式池成员(不删除、不进 pending)——pending 是给未授权
新设备的;它们由扫描线程每轮自动重连,也可前端手动立即重连。
"""
online = _parse_adb_devices(_adb("devices"))
with _ctx():
rows = [d.to_dict() for d in Device.query.filter_by(enabled=True)
.order_by(Device.serial).all()]
return [{"serial": r["serial"], "model": r.get("model") or "",
"name": r.get("name") or ""}
for r in rows if r["serial"] not in online]
# ================== 定时扫描线程 ==================
def _discovery_loop():
"""定时扫描 daemon 线程。每轮重读配置(开关/周期即时生效)。
启动先 sleep 15s 避让 web_server 预连接线程(两者都抢 adb 全局锁)。
等待用 30s 切片(_stop_event.wait),关停/改周期 ≤30s 生效。
"""
try:
_stop_event.wait(15)
while not _stop_event.is_set():
try:
settings = get_settings()
if settings["enabled"]:
scan_once()
except Exception as e:
_log.warning(f"发现定时扫描异常: {e}")
wait_s = max(10, settings.get("interval", 60))
waited = 0
while waited < wait_s and not _stop_event.is_set():
_stop_event.wait(min(30, wait_s - waited))
waited += min(30, wait_s - waited)
except Exception:
pass
def init_app(app):
"""web_server 启动时调用:绑 app + 起定时扫描 daemon 线程。"""
global _app
_app = app
_stop_event.clear()
t = threading.Thread(target=_discovery_loop, name="device-discovery", daemon=True)
t.start()
_log.info("设备自动发现线程已启动(默认 60s 扫描一次)")
def shutdown():
"""优雅退出:置 stop_event,等待中的循环在切片边界退出。"""
_stop_event.set()
# ================== 对外查询与确认 ==================
def get_status():
"""API 状态快照:配置 + 扫描状态 + 最近一轮结果 + 待连接数量。"""
settings = get_settings()
with _ctx():
pending = PendingDevice.query.count()
return {
"enabled": settings["enabled"],
"subnets": settings["subnets"],
"interval": settings["interval"],
"port": settings["port"],
"auto_claim": settings["auto_claim"],
"scanning": _scanning,
"last_scan": _last_scan[0] if _last_scan else "",
"last_result": ({"found": _last_scan[1], "verified": _last_scan[2],
"new": _last_scan[3]} if _last_scan else None),
"last_error": _last_error,
"pending_count": pending,
}
def list_pending():
"""待连接列表:只返回当前在线的设备(离线候选不可确认,不展示)。
后端直接过滤(2026-09-04):不依赖前端 JS 版本,任何客户端都拿不到
离线条目;设备恢复在线后扫描自动更新 last_seen 并重新出现在列表。
"""
online = _parse_adb_devices(_adb("devices"))
with _ctx():
rows = [p.to_dict() for p in PendingDevice.query.order_by(
PendingDevice.first_seen.desc()).all()]
out = []
for r in rows:
if r["serial"] not in online:
continue
r["online"] = True
# 指纹匹配:这台其实就是设备池里某台设备换了地址(前端据此提示"认领")
r["match"] = None
if r.get("fingerprint"):
try:
from core import device_pool
m = device_pool.find_by_fingerprint(r["fingerprint"])
if m and m["serial"] != r["serial"]:
r["match"] = {"serial": m["serial"], "name": m.get("name") or ""}
except Exception:
pass
out.append(r)
return out
def confirm_pending(serial, name="", fingerprint=""):
"""确认连接:pending 行 → 正式设备池(add_device upsert)→ 删 pending。
先按指纹尝试**认领**:同一台物理设备换了地址时,把池中旧记录迁到新 serial,
并同步分组/任务里的引用(名称等信息全部保留),而不是新增一条。
返回 (ok, msg, is_new)。adb_connect + 采型号/指纹由调用方(API 层后台线程)做。
"""
with _ctx():
row = PendingDevice.query.get(serial)
if not row:
return False, "设备不在待连接列表", False
from core import device_pool
fp = (fingerprint or row.fingerprint or "").strip()
claimed_old, claimed_name = device_pool.claim_device(serial, fp)
if claimed_old:
device_pool.add_device(serial, name=claimed_name or name, fingerprint=fp)
PendingDevice.query.filter_by(serial=serial).delete()
db.session.commit()
return True, (f"已认领为『{claimed_name or claimed_old}』"
f"(原地址 {claimed_old},分组/任务的引用已同步)"), False
is_new = device_pool.add_device(serial, name=name or "", fingerprint=fp)
PendingDevice.query.filter_by(serial=serial).delete()
db.session.commit()
return True, "已加入设备池" + ("" if is_new else "(已存在,信息已更新)"), is_new
def ignore_pending(serial):
"""忽略:删除 pending 行(下轮扫描可能再次发现)。"""
with _ctx():
row = PendingDevice.query.get(serial)
if not row:
return False, "设备不在待连接列表"
PendingDevice.query.filter_by(serial=serial).delete()
db.session.commit()
return True, "已忽略"