Files
auto_control/core/apk_manager.py
T

428 lines
18 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
from config import APK_DIR, ADB_PATH
from core.logger import get_logger
from core.models import db, ApkFile as ApkRow
from .stf_client import STFClient
from .device_worker import get_all_worker_status
_log = get_logger("core.apk")
# adb install 超时(秒)。大 APK 安装慢,给 5 分钟
_INSTALL_TIMEOUT = 300
# 并发安装数
_INSTALL_CONCURRENCY = 5
class ApkManager:
"""APK 文件管理 + 批量安装。"""
def __init__(self, stf_client=None, app=None):
self.stf = stf_client or STFClient()
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 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 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, "未选择设备"
# 获取设备信息(型号等)
try:
all_devices = {d["serial"]: d for d in self.stf.list_all_devices()}
except Exception:
all_devices = {}
# 获取 worker 状态(判断设备是否运行中)
worker_status = {w["serial"]: w for w in get_all_worker_status()}
# 构造安装任务
items = {}
for s in serials:
dev = all_devices.get(s, {})
model = dev.get("model") or dev.get("manufacturer") or s
items[s] = {
"name": model,
"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 ThreadPoolExecutor(max_workers=_INSTALL_CONCURRENCY) as pool:
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}")
# 标记完成
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 到设备(大 APK 经 WiFi 推送较慢,超时与安装共用)
self._set_item(serial, "installing", "正在推送 APK...")
r = subprocess.run(
[ADB_PATH, "-s", serial, "push", apk_path, "/data/local/tmp/_install.apk"],
capture_output=True, timeout=_INSTALL_TIMEOUT
)
push_out = ((r.stdout or b"") + (r.stderr or b"")).decode("utf-8", errors="replace")
if r.returncode != 0 or "1 file pushed" not in push_out:
msg = push_out.strip().replace("\n", " ")[:150] or f"推送失败(returncode={r.returncode})"
self._set_item(serial, "failed", msg)
_log.warning(f"[{serial}] APK 推送失败: {msg}")
return
# 3. pm install 强制安装(-r 覆盖安装 / -t 允许测试包 / -d 允许降级版本)
self._set_item(serial, "installing", "正在安装...")
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 _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),
}