Files
auto_control/core/apk_manager.py
T
butubb 7b6c8386af fix(APK 安装): 推送卡死会永久挂住整批安装(新增停滞看门狗 + 收尾兜底)
用户报:220 上一个安装任务卡住不动。(现场:rednote 169.6MB 装 6 台,
5 台成功,192.168.20.203 卡在「推送中 67%(114.0/169.6 MB)」9 分钟没动静——
那台设备当时正跑着任务,adb 流量撞在一起把推送链路卡住了。)

两个叠加的缺陷:
1. 分块推送的 `proc.stdin.write()` **没有任何超时**:链路一卡就永久阻塞,
   安装线程挂死。
2. 批处理的 `as_completed(timeout=...)` 超时后异常直接冒泡出去(`with
   ThreadPoolExecutor` 退出时还会 wait 卡住的线程),`finished=True` 永远没被
   置位 —— 于是界面上永远是"安装中",且后续安装全被「已有安装任务正在进行」
   挡住(用户看到的"卡住")。

修法(core/apk_manager.py):
- 新增**停滞看门狗**:盯着"已写字节数"是否推进,`_PUSH_STALL_TIMEOUT`(90s) 没进展
  就 kill 掉 adb,阻塞中的写立刻以异常返回 → 退回 `adb push` 重传。
  (为什么不用非阻塞写/管道 select:`os.set_blocking` 是 Unix 专有,开发机是 Windows,
  看门狗两边都能用。)
- 批处理改成 try/except/finally:超时或异常时把还没结果的设备标失败,**无论如何**
  都置 `finished=True`;`pool.shutdown(wait=False)` 不再为卡住的线程陪等。

验证:
- 假 adb 模拟"永不读 stdin"(链路卡死)→ 4 秒识别停滞 → kill → 回退 push 成功,
  全程 6.1s(旧代码在这里永久挂死)。
- 真机正常路径不受影响:5.9MB 推送 1.7s、设备侧字节数一致。
2026-09-14 10:28:22 +08:00

693 lines
32 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.
"""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.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)} 台设备")
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}")
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),
}