"""设备自动化管理后台(Flask + Flask-Login,单页应用)。 启动:python web_server.py 访问:http://localhost:5000 页面结构: / — 单页应用(监控/任务/分组/日志/用户 Tab 切换,需登录) /login /logout — 用户登录/登出 /api/* — JSON API 用户系统: 首次启动自动创建默认管理员 admin/admin123(请及时改密码)。 用户数据存 data/users.db(SQLite)。 """ import os import time from flask import Flask, jsonify, request, redirect, url_for, render_template, Response, make_response from flask_login import (LoginManager, login_user, logout_user, login_required, current_user) from core.stf_client import STFClient from core.task_manager import TaskManager from core.logger import get_logger, _LOG_DIR, _MODULE_FILES from core.models import db, init_db, User, DeviceGroup, TaskJob from core.apk_manager import ApkManager from core.adb_helper import screenshot from tasks import list_task_types, get_task_class _log = get_logger("web") app = Flask(__name__) app.config["SECRET_KEY"] = "dev-secret-key-change-in-production" app.config["SQLALCHEMY_DATABASE_URI"] = "sqlite:///" + os.path.join( os.path.dirname(os.path.abspath(__file__)), "data", "users.db") app.config["SQLALCHEMY_TRACK_MODIFICATIONS"] = False login_manager = LoginManager(app) login_manager.login_view = "login" # 先初始化数据库(含旧 JSON 迁移),再创建 TaskManager(需要 app context 读写 DB) init_db(app) stf = STFClient() mgr = TaskManager(stf, app=app) apk_mgr = ApkManager(stf, app=app) @login_manager.user_loader def load_user(user_id): return User.query.get(int(user_id)) # ================== 页面路由 ================== @app.route("/") @login_required def index(): """单页应用首页。""" resp = make_response(render_template("admin/monitor.html")) resp.headers["Cache-Control"] = "no-store, no-cache, must-revalidate, max-age=0" resp.headers["Pragma"] = "no-cache" return resp @app.route("/login", methods=["GET", "POST"]) def login(): if request.method == "POST": username = request.form.get("username", "") password = request.form.get("password", "") user = User.query.filter_by(username=username).first() if user and user.check_password(password): login_user(user) _log.info(f"用户 {username} 登录") return redirect(request.args.get("next") or url_for("index")) return render_template("admin/login.html", error="用户名或密码错误") return render_template("admin/login.html", error=None) @app.route("/logout") @login_required def logout(): _log.info(f"用户 {current_user.username} 登出") logout_user() return redirect(url_for("login")) # ================== API:状态(监控大屏用)================== @app.route("/api/status") @login_required def api_status(): devices, err = mgr.get_status() if err: return jsonify({"ok": False, "error": err}), 500 return jsonify({ "ok": True, "devices": devices, "server_time": time.time(), "fg_scanning": mgr._fg_scanner.is_scanning, "fg_last_scan": mgr._fg_scanner.last_scan_time, }) @app.route("/api/scan_foreground", methods=["POST"]) @login_required def api_scan_foreground(): """手动触发前台 App 扫描(不打扰设备)。""" started = mgr._fg_scanner.scan_once() if started: return jsonify({"ok": True, "msg": "扫描已启动"}) return jsonify({"ok": False, "error": "已有扫描在进行中"}) @app.route("/api/task_types") @login_required def api_task_types(): return jsonify({"ok": True, "task_types": list_task_types()}) @app.route("/api/actions") @login_required def api_actions(): """返回指定任务类型支持的专属操作。""" task_type = request.args.get("task_type", "") task_cls = get_task_class(task_type) if not task_cls: return jsonify({"ok": False, "error": "未知任务类型", "actions": []}) actions = task_cls.list_action_types() return jsonify({"ok": True, "actions": actions}) @app.route("/api/devices") @login_required def api_devices(): """返回所有在线设备 serial(供分组表单勾选用)。""" try: return jsonify({"ok": True, "devices": mgr.list_all_serials()}) except Exception as e: return jsonify({"ok": False, "error": str(e)}), 500 # ================== API:任务计划 CRUD ================== @app.route("/api/jobs") @login_required def api_jobs_list(): return jsonify({"ok": True, "jobs": [j.to_dict() for j in mgr.jobs.values()], "task_types": list_task_types()}) @app.route("/api/jobs", methods=["POST"]) @login_required def api_jobs_create(): data = request.json or {} name = (data.get("name") or "").strip() if not name: return jsonify({"ok": False, "error": "任务名不能为空"}), 400 task_type = data.get("task_type", "douyin_nurture") if not get_task_class(task_type): return jsonify({"ok": False, "error": f"未知任务类型: {task_type}"}), 400 job = mgr.add_job( name=name, task_type=task_type, target=data.get("target", {"mode": "all"}), params=data.get("params", {}), schedule=data.get("schedule", {"mode": "once"}), retry=data.get("retry", {"max_attempts": 1, "delay": 60}), enabled=data.get("enabled", True)) return jsonify({"ok": True, "msg": "任务已创建", "job": job.to_dict()}) @app.route("/api/jobs/", methods=["PUT"]) @login_required def api_jobs_update(job_id): job = mgr.jobs.get(job_id) if not job: return jsonify({"ok": False, "error": "任务不存在"}), 404 data = request.json or {} fields = {} for k in ("name", "task_type", "target", "params", "schedule", "retry", "enabled"): if k in data: fields[k] = data[k] if "task_type" in fields and not get_task_class(fields["task_type"]): return jsonify({"ok": False, "error": f"未知任务类型: {fields['task_type']}"}), 400 mgr.update_job(job_id, **fields) return jsonify({"ok": True, "msg": "任务已更新", "job": job.to_dict()}) @app.route("/api/jobs/", methods=["DELETE"]) @login_required def api_jobs_delete(job_id): if mgr.delete_job(job_id): return jsonify({"ok": True, "msg": "任务已删除"}) return jsonify({"ok": False, "error": "任务不存在"}), 404 @app.route("/api/jobs//run", methods=["POST"]) @login_required def api_jobs_run(job_id): return jsonify(mgr.run_job_now(job_id)) @app.route("/api/jobs//toggle", methods=["POST"]) @login_required def api_jobs_toggle(job_id): enabled = (request.json or {}).get("enabled", True) job = mgr.toggle_job(job_id, enabled) if not job: return jsonify({"ok": False, "error": "任务不存在"}), 404 return jsonify({"ok": True, "msg": f"任务已{'启用' if enabled else '停用'}"}) # ================== API:设备分组 CRUD ================== @app.route("/api/groups") @login_required def api_groups_list(): return jsonify({"ok": True, "groups": [g.to_dict() for g in mgr.groups.values()]}) @app.route("/api/groups", methods=["POST"]) @login_required def api_groups_create(): data = request.json or {} name = (data.get("name") or "").strip() if not name: return jsonify({"ok": False, "error": "分组名不能为空"}), 400 if name in mgr.groups: return jsonify({"ok": False, "error": "分组名已存在"}), 400 mgr.add_group(name, data.get("serials", []), data.get("description", "")) return jsonify({"ok": True, "msg": "分组已创建"}) @app.route("/api/groups/", methods=["PUT"]) @login_required def api_groups_update(name): if name not in mgr.groups: return jsonify({"ok": False, "error": "分组不存在"}), 404 data = request.json or {} mgr.update_group(name, serials=data.get("serials"), description=data.get("description")) return jsonify({"ok": True, "msg": "分组已更新"}) @app.route("/api/groups/", methods=["DELETE"]) @login_required def api_groups_delete(name): if mgr.delete_group(name): return jsonify({"ok": True, "msg": "分组已删除"}) return jsonify({"ok": False, "error": "分组不存在"}), 404 # ================== API:用户管理 CRUD ================== @app.route("/api/users") @login_required def api_users_list(): users = [{"id": u.id, "username": u.username, "is_admin": u.is_admin} for u in User.query.all()] return jsonify({"ok": True, "users": users}) @app.route("/api/users", methods=["POST"]) @login_required def api_users_create(): data = request.json or {} username = (data.get("username") or "").strip() password = data.get("password", "") if not username or not password: return jsonify({"ok": False, "error": "用户名和密码不能为空"}), 400 if User.query.filter_by(username=username).first(): return jsonify({"ok": False, "error": "用户名已存在"}), 400 u = User(username=username, is_admin=data.get("is_admin", True)) u.set_password(password) db.session.add(u) db.session.commit() _log.info(f"创建用户 {username}") return jsonify({"ok": True, "msg": "用户已创建"}) @app.route("/api/users/", methods=["PUT"]) @login_required def api_users_update(uid): u = User.query.get(uid) if not u: return jsonify({"ok": False, "error": "用户不存在"}), 404 data = request.json or {} if "password" in data and data["password"]: u.set_password(data["password"]) if "is_admin" in data: u.is_admin = bool(data["is_admin"]) db.session.commit() _log.info(f"更新用户 {u.username}") return jsonify({"ok": True, "msg": "用户已更新"}) @app.route("/api/users/", methods=["DELETE"]) @login_required def api_users_delete(uid): u = User.query.get(uid) if not u: return jsonify({"ok": False, "error": "用户不存在"}), 404 if u.username == "admin": return jsonify({"ok": False, "error": "不能删除默认管理员"}), 400 if u.id == current_user.id: return jsonify({"ok": False, "error": "不能删除当前登录用户"}), 400 db.session.delete(u) db.session.commit() _log.info(f"删除用户 {u.username}") return jsonify({"ok": True, "msg": "用户已删除"}) # ================== API:日志查看 ================== @app.route("/api/logs") @login_required def api_logs(): files = list(_MODULE_FILES.values()) current = request.args.get("file", "core.log") lines = int(request.args.get("lines", 300)) content = "" path = os.path.join(_LOG_DIR, current) if os.path.exists(path): try: with open(path, encoding="utf-8") as f: content = "".join(f.readlines()[-lines:]) except Exception as e: content = f"读取失败: {e}" return jsonify({"ok": True, "content": content, "file": current, "files": files}) # ================== API:运行控制 ================== @app.route("/api/stop_device", methods=["POST"]) @login_required def api_stop_device(): serial = (request.json or {}).get("serial", "") if mgr.stop_device(serial): return jsonify({"ok": True, "msg": f"已发送停止信号给 {serial}"}) return jsonify({"ok": False, "error": f"{serial} 没有运行中的任务"}), 400 @app.route("/api/stop_all", methods=["POST"]) @login_required def api_stop_all(): stopped = mgr.stop_all() return jsonify({"ok": True, "stopped": stopped}) @app.route("/api/release", methods=["POST"]) @login_required def api_release(): released = stf.release_all_mine() return jsonify({"ok": True, "released": released}) @app.route("/api/device/screenshot") @login_required def api_device_screenshot(): """获取设备当前画面截图(PNG)。 用 adb exec-out screencap -p,只读操作,不抢占 u2 的 atx-agent 通道, 任务运行中调用安全。直接返回 image/png,前端用 加载。 ?serial=xxx 设备 serial ?t=123 时间戳,避免浏览器缓存(前端自动加) """ serial = request.args.get("serial", "") if not serial: return jsonify({"ok": False, "error": "缺少 serial"}), 400 ok, data = screenshot(serial) if ok: return Response(data, mimetype="image/png", headers={"Cache-Control": "no-store"}) return jsonify({"ok": False, "error": data}), 500 # ================== API:应用管理(APK 上传/安装)================== @app.route("/api/apks") @login_required def api_apks_list(): """列出所有已上传的 APK。""" return jsonify({"ok": True, "apks": apk_mgr.list_all()}) @app.route("/api/apks/upload", methods=["POST"]) @login_required def api_apks_upload(): """上传 APK 文件(multipart/form-data,字段名 file)。""" file = request.files.get("file") if not file or not file.filename: return jsonify({"ok": False, "error": "未选择文件"}), 400 info = apk_mgr.upload(file) if info: return jsonify({"ok": True, "apk": info, "msg": f"上传成功: {info['display_name']}"}) return jsonify({"ok": False, "error": "上传失败,请检查文件格式"}), 500 @app.route("/api/apks/", methods=["DELETE"]) @login_required def api_apks_delete(apk_id): """删除 APK 文件和记录。""" ok, msg = apk_mgr.delete(apk_id) if ok: return jsonify({"ok": True, "msg": msg}) return jsonify({"ok": False, "error": msg}), 400 @app.route("/api/apks/install", methods=["POST"]) @login_required def api_apks_install(): """批量安装 APK 到指定设备。参数: {apk_id, serials:[]}""" data = request.json or {} apk_id = data.get("apk_id", "") serials = data.get("serials", []) ok, msg = apk_mgr.install(apk_id, serials) if ok: return jsonify({"ok": True, "msg": msg}) return jsonify({"ok": False, "error": msg}), 400 @app.route("/api/apks/install/status") @login_required def api_apks_install_status(): """获取安装任务实时状态。""" status = apk_mgr.get_install_status() return jsonify({"ok": True, "status": status}) def _run_server(host, port): """启动 Flask,自动处理端口冲突和权限问题。 常见失败原因: WinError 10013 — 端口在 Windows 动态端口范围内被出站连接占用, 或非管理员绑定 0.0.0.0。前者换端口,后者降级 127.0.0.1。 WinError 10048 — 端口已被其他进程监听。换端口重试。 策略:原端口失败 → 试 127.0.0.1:原端口 → 试 127.0.0.1:原端口+1..+5 """ import socket # 先探测原端口是否可用(避免 flask 报错打日志难看) candidates = [(host, port)] if host == "0.0.0.0": candidates.append(("127.0.0.1", port)) for i in range(1, 6): candidates.append(("127.0.0.1", port + i)) for h, p in candidates: try: with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) s.bind((h, p)) # 探测成功,启动 flask if (h, p) != (host, port): _log.warning(f"原端口 {host}:{port} 不可用,改用 {h}:{p}") _log.info(f"启动服务: http://localhost:{p}/") app.run(host=h, port=p, debug=False, threaded=True) return except OSError as e: _log.warning(f"绑定 {h}:{p} 失败: {e}") # 全部失败 raise OSError(f"无法绑定任何候选端口({host}:{port} 及 127.0.0.1:{port}~{port+5})") if __name__ == "__main__": from config import WEB_HOST, WEB_PORT _log.info(f"管理后台: http://localhost:{WEB_PORT}/ (admin/admin123)") try: _run_server(WEB_HOST, WEB_PORT) finally: mgr.shutdown()