"""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 .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, 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 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 到设备(分块 + 实时进度) 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) 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 return False, f"推送中断: {e}" 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) # 校验设备上文件大小一致(exec-in 中途断流可能"成功"但少字节) try: rs = subprocess.run( [ADB_PATH, "-s", serial, "shell", "stat", "-c", "%s", remote], capture_output=True, timeout=30) got = int((rs.stdout or b"").decode("utf-8", "replace").strip().splitlines()[-1]) if got != total: return False, f"推送不完整(设备 {got} / 本地 {total} 字节)" except Exception: pass # stat 参数各机型有差异,取不到就跳过校验 return True, "推送完成 %.1fMB / %.0fs" % (total / 1048576.0, time.time() - t0) 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})") 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), }