平台此前出了问题只能靠人盯页面。现在各组件统一走 `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 功能索引与日志表。
706 lines
33 KiB
Python
706 lines
33 KiB
Python
"""APK 文件管理 + 批量安装。
|
||
|
||
功能:
|
||
- upload: 保存 APK 到 data/apks/,用 pyaxmlparser 解析包名/版本/应用名,入库
|
||
- list_all: 列出所有已上传的 APK
|
||
- delete: 删文件 + 删数据库记录
|
||
- install: 批量安装到指定设备(后台线程,直连设备 adb)
|
||
- get_install_status: 获取安装进度
|
||
|
||
设备连接策略(直连设备 IP:5555,完全绕过 STF occupy/release):
|
||
- worker 运行中:跳过(避免打断任务)
|
||
- 其他设备(自己占用/空闲/他人占用):直接 adb connect serial → install
|
||
不经过 STF,不会触发 STF agent 清理,安装的 app 会永久保留
|
||
- 安装后**不主动 adb disconnect**:serial 是 IP:5555,STF provider 共享该地址
|
||
的 adb transport,disconnect 会让 STF 误判设备离线并触发重连。
|
||
- serial 格式为 IP:5555(设备本身的 adb 网络地址),STF 返回的 serial 即此格式
|
||
|
||
为什么不用 STF occupy/release:
|
||
STF release 会触发 agent 清理设备(卸载第三方 app、清除数据、回桌面),
|
||
导致刚安装的 app 被自动删除。直连设备完全绕过 STF,安装完保持连接即可。
|
||
"""
|
||
import os
|
||
import time
|
||
import uuid
|
||
|
||
import threading
|
||
import subprocess
|
||
from concurrent.futures import (ThreadPoolExecutor, as_completed,
|
||
TimeoutError as FuturesTimeout)
|
||
|
||
from config import APK_DIR, ADB_PATH
|
||
from core.logger import get_logger
|
||
from core import notifier
|
||
from core.models import db, ApkFile as ApkRow
|
||
from .device_worker import get_all_worker_status
|
||
|
||
_log = get_logger("core.apk")
|
||
|
||
# adb install 超时(秒)。大 APK 安装慢,给 5 分钟
|
||
_INSTALL_TIMEOUT = 300
|
||
# 分块推送"多久没有任何字节被接收"就算链路卡死(秒)。注意这是**停滞**判据,
|
||
# 不是总时长:正常推送会不断刷新它。
|
||
_PUSH_STALL_TIMEOUT = 90
|
||
# 并发安装数
|
||
_INSTALL_CONCURRENCY = 5
|
||
|
||
|
||
class ApkManager:
|
||
"""APK 文件管理 + 批量安装。"""
|
||
|
||
def __init__(self, app=None):
|
||
self.app = app # Flask app,用于 db context
|
||
os.makedirs(APK_DIR, exist_ok=True)
|
||
self._lock = threading.Lock()
|
||
self._install_task = None # 当前安装任务
|
||
|
||
def _db(self):
|
||
if self.app is None:
|
||
raise RuntimeError("ApkManager 未关联 Flask app")
|
||
return self.app.app_context()
|
||
|
||
# ================== 上传 ==================
|
||
def upload(self, file_storage):
|
||
"""保存上传的 APK 文件,解析元信息,入库。
|
||
|
||
file_storage: werkzeug FileStorage 对象
|
||
返回: dict (APK 信息) 或 None(失败)
|
||
"""
|
||
# 生成唯一 id
|
||
apk_id = uuid.uuid4().hex[:8]
|
||
original_name = file_storage.filename or f"{apk_id}.apk"
|
||
if not original_name.lower().endswith(".apk"):
|
||
original_name += ".apk"
|
||
disk_filename = f"{apk_id}.apk"
|
||
disk_path = os.path.join(APK_DIR, disk_filename)
|
||
|
||
# 保存文件
|
||
try:
|
||
file_storage.save(disk_path)
|
||
except Exception as e:
|
||
_log.error(f"保存 APK 文件失败: {e}")
|
||
return None
|
||
|
||
size = os.path.getsize(disk_path)
|
||
display_name = os.path.splitext(original_name)[0]
|
||
package_name = ""
|
||
version_name = ""
|
||
version_code = 0
|
||
|
||
# 解析 APK 元信息
|
||
try:
|
||
from pyaxmlparser import APK as APKParser
|
||
apk = APKParser(disk_path)
|
||
package_name = apk.package or ""
|
||
version_name = getattr(apk, "version_name", "") or ""
|
||
version_code = getattr(apk, "version_code", 0) or 0
|
||
# 应用名:优先 get_app_name(),其次 application.label
|
||
try:
|
||
display_name = apk.get_app_name() or display_name
|
||
except Exception:
|
||
try:
|
||
display_name = apk.application.label or display_name
|
||
except Exception:
|
||
pass
|
||
_log.info(f"解析 APK 成功: {display_name} ({package_name} v{version_name})")
|
||
except ImportError:
|
||
_log.warning("pyaxmlparser 未安装,跳过 APK 元信息解析")
|
||
except Exception as e:
|
||
_log.warning(f"解析 APK 元信息失败(不影响上传): {e}")
|
||
|
||
# 入库
|
||
upload_time = time.strftime("%Y-%m-%d %H:%M:%S")
|
||
try:
|
||
with self._db():
|
||
row = ApkRow(id=apk_id, filename=disk_filename,
|
||
display_name=display_name,
|
||
package_name=package_name,
|
||
version_name=version_name,
|
||
version_code=version_code,
|
||
size=size, upload_time=upload_time)
|
||
db.session.add(row)
|
||
db.session.commit()
|
||
_log.info(f"APK 上传成功: {display_name}({apk_id}) {size}字节")
|
||
return row.to_dict()
|
||
except Exception as e:
|
||
_log.error(f"APK 入库失败: {e}")
|
||
# 入库失败但文件已存,仍然返回基本信息
|
||
try:
|
||
os.remove(disk_path)
|
||
except Exception:
|
||
pass
|
||
return None
|
||
|
||
# ================== 列表 ==================
|
||
def list_all(self):
|
||
"""列出所有 APK(**每个上传版本一行**,内部用)。"""
|
||
try:
|
||
with self._db():
|
||
rows = ApkRow.query.order_by(ApkRow.upload_time.desc()).all()
|
||
return [r.to_dict() for r in rows]
|
||
except Exception as e:
|
||
_log.error(f"列出 APK 失败: {e}")
|
||
return []
|
||
|
||
def list_latest(self):
|
||
"""按**包名**合并后的列表:一个应用一行,只显示最新版本。
|
||
|
||
为什么:同一个应用传过多个版本(v1.2 / v1.3)时,列表里会并排出现两行——
|
||
用户看到的是"同一个应用怎么重复了"。这里合并成一行(最新版),
|
||
其余版本的信息放在 `versions` 里,界面可以提示"另有 N 个旧版本"并按需清理。
|
||
|
||
管理页和设备端商店都用这个(设备上选应用时同样不该出现两行同一个 App)。
|
||
解析不出包名的 APK 各自独立成行,不做合并(避免把无关文件混在一起)。
|
||
"""
|
||
rows = self.list_all()
|
||
groups = {}
|
||
for r in rows:
|
||
pkg = (r.get("package_name") or "").strip()
|
||
key = pkg if pkg else ("__single__" + str(r.get("id")))
|
||
groups.setdefault(key, []).append(r)
|
||
|
||
def sort_key(x):
|
||
return (x.get("version_code") or 0, x.get("upload_time") or "")
|
||
|
||
out = []
|
||
for items in groups.values():
|
||
items.sort(key=sort_key, reverse=True)
|
||
top = dict(items[0])
|
||
top["version_count"] = len(items)
|
||
top["versions"] = [{
|
||
"id": i.get("id"),
|
||
"version_name": i.get("version_name") or "",
|
||
"version_code": i.get("version_code") or 0,
|
||
"size": i.get("size") or 0,
|
||
"upload_time": i.get("upload_time") or "",
|
||
} for i in items]
|
||
out.append(top)
|
||
out.sort(key=lambda x: x.get("upload_time") or "", reverse=True)
|
||
return out
|
||
|
||
# ================== 删除 ==================
|
||
def delete(self, apk_id):
|
||
"""删除 APK 文件和数据库记录。"""
|
||
try:
|
||
with self._db():
|
||
row = ApkRow.query.get(apk_id)
|
||
if not row:
|
||
return False, "APK 不存在"
|
||
name = row.display_name or row.filename
|
||
# 删文件
|
||
disk_path = os.path.join(APK_DIR, row.filename)
|
||
if os.path.exists(disk_path):
|
||
try:
|
||
os.remove(disk_path)
|
||
except Exception as e:
|
||
_log.warning(f"删除 APK 文件失败: {e}")
|
||
# 删记录
|
||
db.session.delete(row)
|
||
db.session.commit()
|
||
_log.info(f"已删除 APK: {name}({apk_id})")
|
||
return True, f"已删除 {name}"
|
||
except Exception as e:
|
||
_log.error(f"删除 APK 失败: {e}")
|
||
return False, str(e)
|
||
|
||
def delete_package(self, apk_id, keep_latest=False):
|
||
"""按**包名**删除:一个应用的所有版本一起处理。
|
||
|
||
keep_latest=True → 只删旧版本、保留最新那版("清理旧版本")
|
||
keep_latest=False → 整个应用删掉(所有版本)
|
||
没有包名的 APK 退回单条删除。
|
||
"""
|
||
try:
|
||
with self._db():
|
||
row = ApkRow.query.get(apk_id)
|
||
if not row:
|
||
return False, "APK 不存在"
|
||
pkg = (row.package_name or "").strip()
|
||
if not pkg:
|
||
return self.delete(apk_id)
|
||
|
||
items = ApkRow.query.filter_by(package_name=pkg).all()
|
||
items.sort(key=lambda r: ((r.version_code or 0),
|
||
(r.upload_time or "")), reverse=True)
|
||
targets = items[1:] if keep_latest else items
|
||
if not targets:
|
||
return True, "没有需要清理的旧版本"
|
||
name = row.display_name or pkg
|
||
for r in targets:
|
||
p = os.path.join(APK_DIR, r.filename)
|
||
if os.path.exists(p):
|
||
try:
|
||
os.remove(p)
|
||
except Exception as e:
|
||
_log.warning(f"删除 APK 文件失败: {e}")
|
||
db.session.delete(r)
|
||
db.session.commit()
|
||
_log.info(f"已删除 {name} 的 {len(targets)} 个版本"
|
||
f"({'保留最新' if keep_latest else '整个应用'})")
|
||
if keep_latest:
|
||
return True, f"已清理 {name} 的 {len(targets)} 个旧版本(保留最新)"
|
||
return True, f"已删除 {name}(共 {len(targets)} 个版本)"
|
||
except Exception as e:
|
||
_log.error(f"按包名删除失败: {e}")
|
||
return False, str(e)
|
||
|
||
# ================== 批量安装 ==================
|
||
def install(self, apk_id, serials):
|
||
"""批量安装 APK 到指定设备(后台线程执行)。
|
||
|
||
参数:
|
||
apk_id: APK ID
|
||
serials: 设备 serial 列表
|
||
|
||
返回: (ok, msg)
|
||
"""
|
||
with self._lock:
|
||
if self._install_task and not self._install_task.get("finished"):
|
||
return False, "已有安装任务正在进行,请等待完成"
|
||
|
||
# 查 APK 信息
|
||
try:
|
||
with self._db():
|
||
row = ApkRow.query.get(apk_id)
|
||
if not row:
|
||
return False, "APK 不存在"
|
||
apk_path = os.path.join(APK_DIR, row.filename)
|
||
apk_name = row.display_name or row.filename
|
||
package_name = row.package_name or ""
|
||
if not os.path.exists(apk_path):
|
||
return False, "APK 文件不存在"
|
||
except Exception as e:
|
||
return False, str(e)
|
||
|
||
if not serials:
|
||
return False, "未选择设备"
|
||
|
||
# 设备信息(型号从 worker 状态取,池内无型号信息)
|
||
all_devices = {}
|
||
|
||
# 获取 worker 状态(判断设备是否运行中)
|
||
worker_status = {w["serial"]: w for w in get_all_worker_status()}
|
||
|
||
# 设备名(平台给设备起的名字,如 A08)—— 进度表里显示这个而不是型号:
|
||
# 8 台同型号时型号根本区分不了,名字才行
|
||
pool_names, fp_names = {}, {}
|
||
try:
|
||
from core import device_pool
|
||
# list_devices() 返回字典(含 name/fingerprint)
|
||
for d in device_pool.list_devices():
|
||
if d.get("serial"):
|
||
pool_names[d["serial"]] = d.get("name") or ""
|
||
# USB 设备的 adb serial == ro.serialno == 池里那条记录的指纹
|
||
if d.get("fingerprint"):
|
||
fp_names[d["fingerprint"]] = d.get("name") or ""
|
||
except Exception as e:
|
||
_log.warning(f"读取设备名失败(进度表将显示地址): {e}")
|
||
|
||
# 构造安装任务
|
||
items = {}
|
||
for s in serials:
|
||
items[s] = {
|
||
"name": pool_names.get(s) or fp_names.get(s) or s,
|
||
"status": "pending",
|
||
"msg": "",
|
||
}
|
||
|
||
self._install_task = {
|
||
"apk_id": apk_id,
|
||
"apk_name": apk_name,
|
||
"started_at": time.time(),
|
||
"finished": False,
|
||
"total": len(serials),
|
||
"items": items,
|
||
}
|
||
|
||
# 后台线程执行安装
|
||
t = threading.Thread(target=self._install_worker,
|
||
args=(apk_id, apk_path, apk_name, package_name,
|
||
serials, all_devices, worker_status),
|
||
name="apk-install", daemon=True)
|
||
t.start()
|
||
_log.info(f"开始安装 {apk_name} 到 {len(serials)} 台设备")
|
||
notifier.notify("apk.install.started", apk_id=apk_id, apk_name=apk_name,
|
||
package_name=package_name, total=len(serials))
|
||
return True, f"开始安装 {apk_name} 到 {len(serials)} 台设备"
|
||
|
||
def _install_worker(self, apk_id, apk_path, apk_name, package_name,
|
||
serials, all_devices, worker_status):
|
||
"""后台安装线程:直连设备安装,完全绕过 STF。"""
|
||
# 只跳过 worker 运行中的设备,其他设备(空闲/自己占用/他人占用)都直连安装
|
||
install_list = []
|
||
for s in serials:
|
||
w = worker_status.get(s, {})
|
||
if w.get("status") in ("running", "connecting"):
|
||
self._set_item(s, "skipped", "任务运行中,已跳过")
|
||
continue
|
||
install_list.append(s)
|
||
|
||
_log.info(f"安装任务 {apk_name}: 待安装={len(install_list)}, "
|
||
f"跳过={len(serials)-len(install_list)}")
|
||
|
||
if install_list:
|
||
# 不用 with(它退出时会 wait 卡住的线程):超时后要能立刻收尾,
|
||
# 否则界面永远停在"安装中",后续安装全被"已有安装任务在进行"挡住
|
||
pool = ThreadPoolExecutor(max_workers=_INSTALL_CONCURRENCY)
|
||
try:
|
||
futures = {pool.submit(self._install_one, s, apk_path,
|
||
package_name): s
|
||
for s in install_list}
|
||
for fut in as_completed(futures, timeout=_INSTALL_TIMEOUT + 60):
|
||
s = futures[fut]
|
||
try:
|
||
fut.result()
|
||
except Exception as e:
|
||
self._set_item(s, "failed", f"安装异常: {e}")
|
||
except FuturesTimeout:
|
||
_log.warning(f"安装任务 {apk_name}: 有设备超时未返回,先收尾")
|
||
for s in install_list:
|
||
it = ((self._install_task or {}).get("items") or {}).get(s) or {}
|
||
if it.get("status") in ("pending", "installing"):
|
||
self._set_item(s, "failed", "安装超时(设备无响应)")
|
||
except Exception as e:
|
||
_log.exception(f"安装任务 {apk_name} 异常")
|
||
for s in install_list:
|
||
it = ((self._install_task or {}).get("items") or {}).get(s) or {}
|
||
if it.get("status") in ("pending", "installing"):
|
||
self._set_item(s, "failed", f"安装异常: {e}")
|
||
finally:
|
||
pool.shutdown(wait=False)
|
||
|
||
# 标记完成(**必须**落到这里:没置 finished 的安装任务会一直挂着,
|
||
# 让"已有安装任务正在进行"永久挡住后续安装 —— 2026-09-14 事故)
|
||
if self._install_task:
|
||
self._install_task["finished"] = True
|
||
# 统计
|
||
success = sum(1 for v in self._install_task["items"].values()
|
||
if v["status"] == "success")
|
||
failed = sum(1 for v in self._install_task["items"].values()
|
||
if v["status"] == "failed")
|
||
skipped = sum(1 for v in self._install_task["items"].values()
|
||
if v["status"] == "skipped")
|
||
_log.info(f"安装完成 {apk_name}: 成功={success}, 失败={failed}, 跳过={skipped}")
|
||
# 通知:批量安装结束(带失败设备与原因,运维要知道哪几台没装上)
|
||
try:
|
||
items = ((self._install_task or {}).get("items") or {})
|
||
failed_items = [f"{v.get('name') or s}·{v.get('msg') or ''}"
|
||
for s, v in items.items() if v.get("status") == "failed"][:10]
|
||
notifier.notify("apk.install.finished", apk_name=apk_name,
|
||
success=success, failed=failed, skipped=skipped,
|
||
total=len(items), failed_items=failed_items)
|
||
except Exception as e:
|
||
_log.warning(f"发送安装完成通知失败(不影响安装): {e}")
|
||
|
||
def _install_one(self, serial, apk_path, package_name=""):
|
||
"""直连设备安装 APK。
|
||
|
||
serial 为 adb 网络地址(IP:5555)时直接 adb connect;
|
||
本机 USB 有线设备(serial 无冒号,如 ZY322XXXX)已在 adb 中,无需 connect。
|
||
不经过 STF occupy/release,安装完断开连接,app 永久保留。
|
||
"""
|
||
from .adb_helper import _adb
|
||
self._set_item(serial, "installing", f"正在连接 {serial}...")
|
||
try:
|
||
if ":" in serial:
|
||
# 网络设备:adb connect(带重试,设备网络可能需要几次握手)
|
||
connected = False
|
||
last_out = ""
|
||
for attempt in range(1, 4): # 最多 3 次
|
||
out = _adb("connect", serial)
|
||
last_out = out
|
||
if "connected" in out.lower() and "failed" not in out.lower():
|
||
connected = True
|
||
break
|
||
_log.info(f"[{serial}] adb connect 第{attempt}次: {out}")
|
||
time.sleep(1 + attempt)
|
||
if not connected:
|
||
_log.warning(f"[{serial}] adb connect 失败: 输出={last_out}")
|
||
self._set_item(serial, "failed", f"adb connect 失败: {last_out[:150]}")
|
||
return
|
||
else:
|
||
# 本机 USB 设备:确认已在 adb devices 列表中
|
||
dev_out = _adb("devices")
|
||
if serial not in (dev_out or ""):
|
||
self._set_item(serial, "failed",
|
||
"设备未在本地 adb 列表中,请检查 USB 连接/驱动")
|
||
return
|
||
|
||
# 强制安装:绕过厂商安全防护(MIUI 等用 adb install 会弹"通过 USB 安装应用"
|
||
# 确认框需要手点;改为 推送 + pm install,shell 安装不弹窗)
|
||
# 1. 关闭 adb 安装校验(包校验器/安全守护拦截)
|
||
try:
|
||
_adb("-s", serial, "shell", "settings", "put", "global",
|
||
"verifier_verify_adb_installs", "0")
|
||
except Exception:
|
||
pass
|
||
# 2. 推送 APK 到设备(分块 + 实时进度)
|
||
ok_push, msg_push = self._push_with_progress(serial, apk_path)
|
||
if not ok_push:
|
||
self._set_item(serial, "failed", msg_push)
|
||
_log.warning(f"[{serial}] APK 推送失败: {msg_push}")
|
||
return
|
||
# 3. pm install 强制安装(-r 覆盖安装 / -t 允许测试包 / -d 允许降级版本)
|
||
# 这一步在**设备端**解包/优化,耗时随包大小增长(抖音 336MB 实测 ~65s),
|
||
# 但系统没给进度接口,只能给出预期时间,别让用户以为是卡住了。
|
||
self._set_item(serial, "installing",
|
||
"正在设备上安装...(大包要 1-2 分钟,无进度可读)")
|
||
r = subprocess.run(
|
||
[ADB_PATH, "-s", serial, "shell", "pm", "install", "-r", "-t", "-d",
|
||
"/data/local/tmp/_install.apk"],
|
||
capture_output=True, timeout=_INSTALL_TIMEOUT
|
||
)
|
||
# 4. 清理设备临时文件(失败残留下次推送会覆盖,不影响)
|
||
try:
|
||
_adb("-s", serial, "shell", "rm", "-f", "/data/local/tmp/_install.apk")
|
||
except Exception:
|
||
pass
|
||
out = ((r.stdout or b"") + (r.stderr or b""))
|
||
try:
|
||
out = out.decode("utf-8", errors="replace")
|
||
except Exception:
|
||
out = out.decode("gbk", errors="replace")
|
||
out = out.strip()
|
||
_log.info(f"[{serial}] adb install 返回码={r.returncode}, 输出={out[:300]}")
|
||
low = out.lower()
|
||
install_ok = False
|
||
if "success" in low:
|
||
install_ok = True
|
||
elif "already_installed" in low or "already installed" in low:
|
||
install_ok = True
|
||
|
||
if not install_ok:
|
||
msg = out.replace("\n", " ").strip()[:200]
|
||
self._set_item(serial, "failed", msg or "安装失败(returncode=%d)" % r.returncode)
|
||
_log.warning(f"[{serial}] APK 安装失败: {msg}")
|
||
return
|
||
|
||
# 安装后验证:查 pm list packages 确认包名是否真正存在
|
||
if package_name:
|
||
self._set_item(serial, "installing", "正在验证安装结果...")
|
||
try:
|
||
rv = subprocess.run(
|
||
[ADB_PATH, "-s", serial, "shell", "pm", "list", "packages", package_name],
|
||
capture_output=True, timeout=30
|
||
)
|
||
vout = (rv.stdout or b"").decode("utf-8", errors="replace").strip()
|
||
_log.info(f"[{serial}] pm list packages {package_name}: {vout[:200]}")
|
||
if package_name in vout:
|
||
self._set_item(serial, "success", "安装成功(已验证)")
|
||
_log.info(f"[{serial}] APK 安装成功且已验证")
|
||
else:
|
||
self._set_item(serial, "failed", "adb返回Success但设备上未检测到包,可能被系统拦截")
|
||
_log.warning(f"[{serial}] 安装验证失败: adb返回Success但pm list未找到 {package_name}")
|
||
except Exception as e:
|
||
self._set_item(serial, "success", "安装成功(验证异常)")
|
||
_log.warning(f"[{serial}] 安装验证异常: {e}")
|
||
else:
|
||
self._set_item(serial, "success", "安装成功")
|
||
_log.info(f"[{serial}] APK 安装成功(无包名,跳过验证)")
|
||
except subprocess.TimeoutExpired:
|
||
self._set_item(serial, "failed", "安装超时")
|
||
except Exception as e:
|
||
self._set_item(serial, "failed", str(e)[:200])
|
||
_log.exception(f"[{serial}] 安装异常")
|
||
finally:
|
||
# 不主动 adb disconnect:serial 是 IP:5555,STF provider 共享该地址的
|
||
# adb transport,disconnect 会让 STF 误判设备离线并触发重连。
|
||
# adb 连接保持即可,不影响设备和后续操作。
|
||
pass
|
||
|
||
def _push_with_progress(self, serial, apk_path,
|
||
remote="/data/local/tmp/_install.apk"):
|
||
"""把 APK 分块写进设备,边传边更新进度。返回 (ok, msg)。
|
||
|
||
为什么不用 `adb push`:它只在**结束时**吐一行汇总
|
||
("1 file pushed, 3.0 MB/s"),大包(抖音 336MB ≈ 1.8 分钟/台)整个传输过程
|
||
界面上只有一句"正在推送 APK...",用户只能干等 —— 这就是要优化的问题。
|
||
|
||
改用 `exec-in` 分块写:每块回调一次进度(百分比/已传/速度/剩余时间)。
|
||
实测速度略降(2.6 vs 3.1 MB/s),换来看得见的进度,划算。
|
||
传完校验设备上的文件大小;exec-in 不可用时退回原来的 `adb push`。
|
||
"""
|
||
total = os.path.getsize(apk_path)
|
||
chunk = 1 << 20 # 1MB:再大提升有限,再小 adb 调用开销明显
|
||
t0 = time.time()
|
||
sent = 0
|
||
last_report = 0.0
|
||
try:
|
||
proc = subprocess.Popen(
|
||
[ADB_PATH, "-s", serial, "exec-in", "cat > %s" % remote],
|
||
stdin=subprocess.PIPE, stdout=subprocess.DEVNULL,
|
||
stderr=subprocess.PIPE)
|
||
except Exception as e:
|
||
_log.warning(f"[{serial}] exec-in 启动失败,退回 adb push: {e}")
|
||
return self._push_fallback(serial, apk_path, remote)
|
||
|
||
# 看门狗:链路卡死时 `proc.stdin.write()` 会**永久阻塞**(设备掉线、
|
||
# WiFi 断了、任务在抢 adb…)→ 线程挂死 → 整批安装永不 finished → 界面
|
||
# 卡在"安装中"、后续安装全被"已有安装任务在进行"挡住(2026-09-14 事故)。
|
||
# 这里盯住"已写字节数"是否还在推进,停滞超时就 kill 掉 adb —— 阻塞中的
|
||
# write 会立刻以异常退出,走下面的回退逻辑。
|
||
# 为什么用看门狗而不是非阻塞写:`os.set_blocking`/管道 select 是 Unix 专有,
|
||
# 开发机是 Windows,看门狗两边都能用。
|
||
stop_watch = threading.Event()
|
||
|
||
def _stall_watchdog():
|
||
last_seen, since = -1, time.time()
|
||
while not stop_watch.wait(2.0):
|
||
if sent == total:
|
||
return
|
||
if sent != last_seen:
|
||
last_seen, since = sent, time.time()
|
||
elif time.time() - since > _PUSH_STALL_TIMEOUT:
|
||
_log.warning(f"[{serial}] 推送停滞 {int(time.time() - since)} 秒"
|
||
f"({sent}/{total} 字节),杀掉推送进程并回退")
|
||
try:
|
||
proc.kill()
|
||
except Exception:
|
||
pass
|
||
return
|
||
|
||
watchdog = threading.Thread(target=_stall_watchdog, daemon=True)
|
||
watchdog.start()
|
||
try:
|
||
with open(apk_path, "rb") as f:
|
||
while True:
|
||
buf = f.read(chunk)
|
||
if not buf:
|
||
break
|
||
proc.stdin.write(buf)
|
||
sent += len(buf)
|
||
now = time.time()
|
||
if now - last_report >= 0.8 or sent >= total:
|
||
last_report = now
|
||
spent = max(now - t0, 0.001)
|
||
speed = sent / spent / 1048576.0 # MB/s
|
||
left = int((total - sent) / max(sent / spent, 1))
|
||
self._set_item(
|
||
serial, "installing",
|
||
"推送中 %d%%(%.1f/%.1f MB · %.1f MB/s · 剩约 %ds)"
|
||
% (sent * 100 // max(total, 1), sent / 1048576.0,
|
||
total / 1048576.0, speed, left))
|
||
proc.stdin.close()
|
||
except Exception as e:
|
||
try:
|
||
proc.kill()
|
||
except Exception:
|
||
pass
|
||
_log.warning(f"[{serial}] exec-in 推送中断({e}),已传 {sent}/{total} 字节,"
|
||
f"退回 adb push 重传")
|
||
return self._push_fallback(serial, apk_path, remote)
|
||
finally:
|
||
stop_watch.set()
|
||
|
||
try:
|
||
rc = proc.wait(timeout=_INSTALL_TIMEOUT)
|
||
except subprocess.TimeoutExpired:
|
||
proc.kill()
|
||
return False, "推送超时(设备无响应)"
|
||
if rc != 0:
|
||
err = ""
|
||
try:
|
||
err = (proc.stderr.read() or b"").decode("utf-8", "replace")[:150]
|
||
except Exception:
|
||
pass
|
||
_log.warning(f"[{serial}] exec-in 推送到一半失败(rc={rc}),退回 adb push")
|
||
return self._push_fallback(serial, apk_path, remote)
|
||
|
||
# 校验设备上文件大小。**必须轮询等它长齐**——见 _wait_remote_size 的说明:
|
||
# exec-in 返回后设备侧可能仍在落盘,只量一次会把成功的传输误判成截断。
|
||
got = self._wait_remote_size(serial, remote, total)
|
||
if got is not None and got != total:
|
||
_log.warning(f"[{serial}] exec-in 推送不完整(设备 {got} / 本地 {total} 字节),"
|
||
f"改用 adb push 重传")
|
||
return self._push_fallback(serial, apk_path, remote)
|
||
|
||
return True, "推送完成 %.1fMB / %.0fs" % (total / 1048576.0, time.time() - t0)
|
||
|
||
@staticmethod
|
||
def _remote_size(serial, remote):
|
||
"""设备上文件字节数;取不到返回 None(stat 参数各机型有差异,不据此判失败)。"""
|
||
try:
|
||
rs = subprocess.run(
|
||
[ADB_PATH, "-s", serial, "shell", "stat", "-c", "%s", remote],
|
||
capture_output=True, timeout=30)
|
||
lines = (rs.stdout or b"").decode("utf-8", "replace").strip().splitlines()
|
||
return int(lines[-1]) if lines else None
|
||
except Exception:
|
||
return None
|
||
|
||
def _wait_remote_size(self, serial, remote, want, timeout=20.0):
|
||
"""轮询设备上的文件大小,直到等于 want / 超时 / 取不到(None)。
|
||
|
||
为什么不能只量一次:**adb 进程退出 ≠ 设备上文件写完**。exec-in 把数据交给
|
||
设备侧(adbd → shell → `cat >`)后,缓冲区里的数据还会继续落盘。实测同一份
|
||
5.9MB 的 APK:adb 刚返回时 stat 读到 3710967,紧接着 4264448,几秒后才是完整的
|
||
5933150。旧实现量一次就判"推送不完整",把**成功的传输误报成失败**
|
||
(2026-09-14 用户报的就是这个;同一份文件在多台设备上"失败",断点还各不相同)。
|
||
"""
|
||
deadline = time.time() + timeout
|
||
last = None
|
||
while True:
|
||
last = self._remote_size(serial, remote)
|
||
if last == want or last is None:
|
||
return last
|
||
if time.time() >= deadline:
|
||
return last
|
||
time.sleep(0.3)
|
||
|
||
def _push_fallback(self, serial, apk_path, remote):
|
||
"""老 adb 或 exec-in 失败时的退路:原来的 `adb push`(无进度)。"""
|
||
self._set_item(serial, "installing", "正在推送 APK...(无进度显示)")
|
||
try:
|
||
r = subprocess.run(
|
||
[ADB_PATH, "-s", serial, "push", apk_path, remote],
|
||
capture_output=True, timeout=_INSTALL_TIMEOUT)
|
||
except subprocess.TimeoutExpired:
|
||
return False, "推送超时"
|
||
out = ((r.stdout or b"") + (r.stderr or b"")).decode("utf-8", errors="replace")
|
||
if r.returncode != 0 or "1 file pushed" not in out:
|
||
return False, (out.strip().replace("\n", " ")[:150]
|
||
or f"推送失败(returncode={r.returncode})")
|
||
# adb push 是同步的(不会出现 exec-in 那种"退出后才落盘"),但仍量一次大小:
|
||
# 设备空间不足/文件被并发覆盖时,"1 file pushed" 也可能是残缺的
|
||
total = os.path.getsize(apk_path)
|
||
got = self._wait_remote_size(serial, remote, total, timeout=10)
|
||
if got is not None and got != total:
|
||
return False, f"推送不完整(设备 {got} / 本地 {total} 字节)"
|
||
return True, "推送完成(adb push)"
|
||
|
||
def _set_item(self, serial, status, msg=""):
|
||
"""更新安装任务中某设备的状态。"""
|
||
if not self._install_task:
|
||
return
|
||
items = self._install_task["items"]
|
||
if serial not in items:
|
||
items[serial] = {"name": serial, "status": status, "msg": msg}
|
||
else:
|
||
items[serial]["status"] = status
|
||
items[serial]["msg"] = msg
|
||
|
||
def get_install_status(self):
|
||
"""获取当前安装任务状态。"""
|
||
if not self._install_task:
|
||
return None
|
||
task = self._install_task
|
||
items = task["items"]
|
||
success = sum(1 for v in items.values() if v["status"] == "success")
|
||
failed = sum(1 for v in items.values() if v["status"] == "failed")
|
||
skipped = sum(1 for v in items.values() if v["status"] == "skipped")
|
||
installing = sum(1 for v in items.values() if v["status"] == "installing")
|
||
pending = sum(1 for v in items.values() if v["status"] == "pending")
|
||
return {
|
||
"apk_id": task["apk_id"],
|
||
"apk_name": task["apk_name"],
|
||
"started_at": task["started_at"],
|
||
"finished": task.get("finished", False),
|
||
"total": task["total"],
|
||
"success": success,
|
||
"failed": failed,
|
||
"skipped": skipped,
|
||
"installing": installing,
|
||
"pending": pending,
|
||
"items": dict(items),
|
||
}
|