原来建表有三套来源(ORM create_all + 自建 SCHEMA_MIGRATIONS + agent_api 里的裸
CREATE TABLE)。这在 SQLite 下能活是因为 SQLite 的 DDL 极宽容;换 MySQL 后
agent_* 四张表的 `INTEGER PRIMARY KEY AUTOINCREMENT`、`TEXT DEFAULT ''`、
`TEXT PRIMARY KEY` 会全部建不出来,而失败被 `except: pass` 吞掉 —— 表现为
「经验库/动作库/会话功能静默失灵、日志里什么都看不到」。本期把它收敛成一套。
- core/models.py:
* 新增 5 个 ORM 模型:AppMeta / AgentExperience / ExperienceAudit /
AgentAction / AgentConversation(列名沿用历史,`app_meta.key` 在 MySQL 里是
保留字,ORM 属性名用 k,读写统一走 db_config.meta_get/meta_set)
* 新增 _long_text()(MySQL 用 MEDIUMTEXT,裸 TEXT 只有 64KB)与 _DOUBLE()
(评分别用单精度 FLOAT);_migrate_schema 不再执行 DDL,只维护版本账本,
SCHEMA_MIGRATIONS 的历史建表/加列条目 SQL 置 None
* 新增 _sync_columns():用 sqlalchemy.inspect 比对模型与实表补缺列,取代原来
靠 "duplicate column name" 报错文本判断的写法(换方言就失效)
* _ensure_unique_indexes() 按方言分叉:MySQL 5.7 没有过滤索引,改用
「虚拟生成列 + 唯一索引」复刻"空值不参与唯一约束"的语义
- web/agent_api.py: 删 4 段裸 CREATE TABLE;_ensure_* 收敛为 _ensure_tables()
(db.create_all 薄封装,失败必记日志);app_meta 读写与 INSERT OR IGNORE
改为方言中立
- core/device_discovery.py: app_meta 读写改用同一助手
- core/system_backup.py: SUMMARY_TABLES 由 db.metadata 派生 —— 备份覆盖红线
从"靠人记"变成结构上不可能漏;CURRENT_SCHEMA_VERSION 改从 models 导入
验证(均在 users.db 的一致快照副本上,不碰真实库):
空库 → 自建 12 张表 + 默认管理员 + 唯一索引 + schema_version=6;
老库 → 逐表行数与真实库完全一致,app_meta 键未变;
接口 → health/devices/jobs/groups/agent(经验库/动作库/会话) 全 200,
会话增删与备份导出(12 张表、无覆盖缺失)正常。
426 lines
17 KiB
Python
426 lines
17 KiB
Python
"""系统数据备份导出/导入(数据库 + APK 文件)。
|
||
|
||
导出:用 sqlite3 在线备份 API 对 data/users.db 做一致性快照 → zip
|
||
(users.db + manifest.json + 可选 apks/*.apk)。
|
||
导入:上传(zip/db) → 暂存校验预览 → 应用(先把当前库快照到
|
||
BACKUP_DIR/pre_restore_*.db,再把暂存库移到 RESTORE_PENDING_DIR)→
|
||
下次启动 consume_pending_restore() 换位生效。
|
||
|
||
为什么必须重启生效:web_server 单进程内 SQLAlchemy engine + 多个常驻后台线程
|
||
(TaskManager 内存态 / 设备发现 / Agent 等)持有 users.db(WAL);Windows 下
|
||
不能直接替换正被打开的文件,且 TaskManager 启动时才把 groups/jobs 读入内存。
|
||
故导入只落一个「待生效恢复任务」,由 web_server.py 在 init_db 之前消费。
|
||
|
||
安全红线:全程只读现有库 / 新增或移动自己目录下的文件;绝不 kill-server、
|
||
绝不 disconnect(与全项目一致)。
|
||
"""
|
||
import json
|
||
import os
|
||
import shutil
|
||
import sqlite3
|
||
import time
|
||
import uuid
|
||
import zipfile
|
||
from urllib.parse import quote
|
||
|
||
from config import (DATA_DIR, APK_DIR, BACKUP_DIR,
|
||
RESTORE_STAGING_DIR, RESTORE_PENDING_DIR)
|
||
from core.logger import get_logger
|
||
from core.models import CURRENT_SCHEMA_VERSION, db
|
||
|
||
_log = get_logger("core.backup")
|
||
|
||
DB_FILE = os.path.join(DATA_DIR, "users.db")
|
||
|
||
# 判定「本平台备份库」的必需表(缺任何一张即拒绝导入)
|
||
REQUIRED_TABLES = ("app_meta", "user", "task_job", "device_group")
|
||
# 备份覆盖清单:**由模型元数据派生**,不再手工维护。
|
||
# 这样「新增持久化表必须登记进备份清单」这条红线从"靠人记"变成结构上不可能漏
|
||
# (2026-09-10 动作库 agent_action 就漏过:数据其实在快照里,只是清单没列,
|
||
# 导出预览里看不到那一行 → 被误判成"没备份")。
|
||
SUMMARY_TABLES = tuple(sorted(t.name for t in db.metadata.tables.values()))
|
||
TABLE_LABELS = {
|
||
"app_meta": "系统配置(app_meta)",
|
||
"user": "用户", "device_group": "设备分组", "task_job": "任务计划",
|
||
"custom_action": "自定义动作", "apk_file": "APK 记录", "device": "设备池",
|
||
"pending_device": "待连接设备", "agent_conversation": "AI 会话",
|
||
"agent_experience": "经验库", "experience_audit": "经验巡检",
|
||
"agent_action": "动作库",
|
||
}
|
||
_STAGE_TTL = 1800 # 导入暂存有效期(秒)
|
||
|
||
|
||
class BackupError(Exception):
|
||
"""备份/导入相关可预期错误(msg 直接给前端展示)。"""
|
||
|
||
|
||
# ================== 通用工具 ==================
|
||
def _ts():
|
||
return time.strftime("%Y%m%d_%H%M%S")
|
||
|
||
|
||
def _readonly_uri(path):
|
||
"""把文件路径转成 sqlite 只读 URI(兼容含空格/中文/反斜杠的 Windows 路径)。"""
|
||
return "file:{}?mode=ro".format(quote(os.path.abspath(path).replace("\\", "/"), safe="/:"))
|
||
|
||
|
||
def _connect_readonly(path):
|
||
con = sqlite3.connect(_readonly_uri(path), uri=True)
|
||
con.text_factory = str # 与平台一致按 UTF-8 读(表内文本均为 UTF-8)
|
||
return con
|
||
|
||
|
||
def remove_quiet(path):
|
||
"""尽力删除文件/目录(发送完成清理等场景,忽略不存在)。"""
|
||
try:
|
||
if os.path.isdir(path):
|
||
shutil.rmtree(path, ignore_errors=True)
|
||
elif os.path.exists(path):
|
||
os.remove(path)
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
def _ensure_dir(path):
|
||
os.makedirs(path, exist_ok=True)
|
||
|
||
|
||
def _list_tables(con):
|
||
rows = con.execute(
|
||
"SELECT name FROM sqlite_master WHERE type='table'").fetchall()
|
||
return {r[0] for r in rows}
|
||
|
||
|
||
def _read_schema_version(con):
|
||
try:
|
||
row = con.execute(
|
||
"SELECT value FROM app_meta WHERE key='schema_version'").fetchone()
|
||
return int(row[0]) if row and row[0] else 0
|
||
except Exception:
|
||
return 0
|
||
|
||
|
||
def _table_rows(con, table):
|
||
try:
|
||
return con.execute('SELECT COUNT(*) FROM "%s"' % table).fetchone()[0]
|
||
except Exception:
|
||
return None
|
||
|
||
|
||
# ================== 在线快照 ==================
|
||
def snapshot_db(dest_path, src_path=None):
|
||
"""对 src(默认当前 users.db)做 sqlite 在线一致快照到 dest_path。
|
||
|
||
用标准库备份 API(不依赖 SQLAlchemy engine,WAL 模式下读一致快照安全)。
|
||
备份完成后把目标强制落回单文件(journal_mode=DELETE + checkpoint),
|
||
避免副本以 WAL 模式残留 -wal/-shm 或主文件缺刚提交帧。
|
||
"""
|
||
src_path = src_path or DB_FILE
|
||
if not os.path.exists(src_path):
|
||
raise BackupError(f"数据库不存在: {src_path}")
|
||
con = _connect_readonly(src_path)
|
||
out = sqlite3.connect(dest_path)
|
||
try:
|
||
con.backup(out)
|
||
try:
|
||
out.execute("PRAGMA journal_mode=DELETE") # checkpoint 并转回 DELETE
|
||
except sqlite3.Error:
|
||
pass
|
||
finally:
|
||
out.close()
|
||
con.close()
|
||
# 防御:确保没有任何残留附属文件
|
||
remove_quiet(dest_path + "-wal")
|
||
remove_quiet(dest_path + "-shm")
|
||
|
||
|
||
# ================== 导出 ==================
|
||
def _summary_info(con):
|
||
"""表行数摘要(只列出存在的表)。"""
|
||
rows = []
|
||
for t in SUMMARY_TABLES:
|
||
n = _table_rows(con, t)
|
||
if n is not None:
|
||
rows.append({"table": t, "label": TABLE_LABELS.get(t, t), "rows": n})
|
||
return rows
|
||
|
||
|
||
def _prune_old_exports(max_age=3600):
|
||
"""清理过期的导出临时 zip(下载完成后的 call_on_close 在 Windows 上可能
|
||
因文件锁删不掉,这里按时间兜底清理;保留近 1 小时的便于失败重试)。"""
|
||
now = time.time()
|
||
if not os.path.isdir(BACKUP_DIR):
|
||
return
|
||
for n in os.listdir(BACKUP_DIR):
|
||
if n.startswith("export_") and n.endswith(".zip"):
|
||
p = os.path.join(BACKUP_DIR, n)
|
||
try:
|
||
if now - os.path.getmtime(p) > max_age:
|
||
os.remove(p)
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
def create_export(include_apk=True):
|
||
"""生成导出 zip,返回 (zip_path, filename, manifest)。"""
|
||
_ensure_dir(BACKUP_DIR)
|
||
_prune_old_exports()
|
||
base = os.path.join(BACKUP_DIR, f"export_{_ts()}")
|
||
snap = base + ".snapshot.db"
|
||
snapshot_db(snap)
|
||
|
||
apk_meta = []
|
||
if include_apk and os.path.isdir(APK_DIR):
|
||
for fname in sorted(os.listdir(APK_DIR)):
|
||
if fname.lower().endswith(".apk"):
|
||
fp = os.path.join(APK_DIR, fname)
|
||
try:
|
||
apk_meta.append({"name": fname, "size": os.path.getsize(fp)})
|
||
except OSError:
|
||
pass
|
||
|
||
zip_path = base + ".zip"
|
||
con = _connect_readonly(snap)
|
||
try:
|
||
manifest = {
|
||
"format": "auto_control_backup",
|
||
"version": 1,
|
||
"created_at": time.strftime("%Y-%m-%d %H:%M:%S"),
|
||
"schema_version": _read_schema_version(con),
|
||
"include_apk": include_apk,
|
||
"tables": _summary_info(con),
|
||
"apks": apk_meta,
|
||
}
|
||
# 覆盖自检:登记在册的业务表若在快照里缺失(新增功能忘了登记 / 建表失败),
|
||
# 显式告警并写进 manifest——避免"以为备份了其实没有"(备份覆盖红线)。
|
||
_present = {t["table"] for t in manifest["tables"]}
|
||
_missing = [t for t in SUMMARY_TABLES if t not in _present]
|
||
if _missing:
|
||
manifest["coverage_missing"] = _missing
|
||
_log.warning(f"备份覆盖检查:以下登记表未纳入快照,请确认是否应备份: {_missing}")
|
||
finally:
|
||
con.close()
|
||
|
||
try:
|
||
with zipfile.ZipFile(zip_path, "w", zipfile.ZIP_DEFLATED) as z:
|
||
z.write(snap, arcname="users.db")
|
||
for a in apk_meta:
|
||
z.write(os.path.join(APK_DIR, a["name"]), arcname="apks/" + a["name"])
|
||
z.writestr("manifest.json", json.dumps(
|
||
manifest, ensure_ascii=False, indent=2))
|
||
finally:
|
||
remove_quiet(snap)
|
||
return zip_path, os.path.basename(zip_path), manifest
|
||
|
||
|
||
# ================== 校验 ==================
|
||
def validate_backup(db_file):
|
||
"""校验备份库可读、是本平台库。返回 (ok, info|error_msg)。"""
|
||
if not os.path.exists(db_file):
|
||
return False, "缺少数据库文件 users.db"
|
||
try:
|
||
con = _connect_readonly(db_file)
|
||
except sqlite3.Error as e:
|
||
return False, f"无法打开数据库: {e}"
|
||
try:
|
||
integrity = con.execute("PRAGMA integrity_check").fetchone()[0]
|
||
if integrity != "ok":
|
||
return False, f"数据库完整性校验失败: {integrity}"
|
||
tables = _list_tables(con)
|
||
missing = [t for t in REQUIRED_TABLES if t not in tables]
|
||
if missing:
|
||
return False, f"不是本平台备份库(缺少必需表: {', '.join(missing)})"
|
||
schema_version = _read_schema_version(con)
|
||
missing_optional = [t for t in SUMMARY_TABLES
|
||
if t not in REQUIRED_TABLES and t not in tables]
|
||
warnings = []
|
||
if schema_version < CURRENT_SCHEMA_VERSION:
|
||
warnings.append(
|
||
f"备份 schema 较旧(v{schema_version} < 当前 v{CURRENT_SCHEMA_VERSION}),"
|
||
f"应用后首次启动会自动迁移补齐结构与新表")
|
||
elif schema_version > CURRENT_SCHEMA_VERSION:
|
||
warnings.append(
|
||
f"备份 schema 较新(v{schema_version} > 当前 v{CURRENT_SCHEMA_VERSION}),"
|
||
f"当前平台版本可能读不了新增列/表,建议先升级再恢复")
|
||
if missing_optional:
|
||
warnings.append("备份缺少部分可选表(" + "、".join(
|
||
TABLE_LABELS.get(t, t) for t in missing_optional)
|
||
+ "),应用后启动会自动补建空表")
|
||
# 反向自检:备份里有"未登记"的表 → 提醒把它纳入覆盖清单(红线)
|
||
extra_tables = sorted(t for t in tables
|
||
if t not in SUMMARY_TABLES and not t.startswith("sqlite_"))
|
||
if extra_tables:
|
||
warnings.append("备份含未登记的其它表(" + "、".join(extra_tables)
|
||
+ ")——如属业务数据,请登记进 core/system_backup.py 的覆盖清单")
|
||
warnings.append("备份为全量快照:含用户口令哈希、AI 配置里的 API Key 等敏感信息,请妥善保管")
|
||
return True, {
|
||
"integrity": integrity,
|
||
"schema_version": schema_version,
|
||
"current_schema_version": CURRENT_SCHEMA_VERSION,
|
||
"tables": _summary_info(con),
|
||
"missing_optional": missing_optional,
|
||
"extra_tables": extra_tables,
|
||
"warnings": warnings,
|
||
}
|
||
except sqlite3.Error as e:
|
||
return False, f"读取数据库失败: {e}"
|
||
finally:
|
||
con.close()
|
||
|
||
|
||
# ================== 导入:暂存 + 预览 ==================
|
||
def _extract_zip(zip_path, staging_dir):
|
||
"""解压备份 zip:取顶层 users.db 与可选 apks/*.apk(防路径穿越,只用 basename)。"""
|
||
db_dest = os.path.join(staging_dir, "users.db")
|
||
apk_dest = os.path.join(staging_dir, "apks")
|
||
got_db = False
|
||
with zipfile.ZipFile(zip_path) as z:
|
||
for info in z.infolist():
|
||
if info.is_dir():
|
||
continue
|
||
name = info.filename.replace("\\", "/")
|
||
base = os.path.basename(name)
|
||
if name == "users.db" and base == "users.db":
|
||
z.extract(info, staging_dir) # 已确认顶层名字,无穿越
|
||
got_db = True
|
||
elif name.startswith("apks/") and base.lower().endswith(".apk"):
|
||
_ensure_dir(apk_dest)
|
||
with z.open(info) as src, open(
|
||
os.path.join(apk_dest, base), "wb") as dst:
|
||
shutil.copyfileobj(src, dst)
|
||
if not got_db:
|
||
raise BackupError("zip 内未找到 users.db(顶层)")
|
||
return db_dest
|
||
|
||
|
||
def _prune_stale_staging():
|
||
"""清理超时未应用的暂存目录(TTL 后自动删除)。"""
|
||
now = time.time()
|
||
if not os.path.isdir(RESTORE_STAGING_DIR):
|
||
return
|
||
for name in os.listdir(RESTORE_STAGING_DIR):
|
||
p = os.path.join(RESTORE_STAGING_DIR, name)
|
||
try:
|
||
if os.path.isdir(p) and now - os.path.getmtime(p) > _STAGE_TTL:
|
||
remove_quiet(p)
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
def stage_upload(file_storage):
|
||
"""保存上传备份并校验,返回 (token, info)。失败抛 BackupError。"""
|
||
_prune_stale_staging()
|
||
orig_name = file_storage.filename or "backup"
|
||
token = uuid.uuid4().hex[:12]
|
||
staging_dir = os.path.join(RESTORE_STAGING_DIR, token)
|
||
_ensure_dir(staging_dir)
|
||
low = (orig_name or "").lower()
|
||
try:
|
||
raw_path = os.path.join(staging_dir, "upload" + ("." + low.rsplit(".", 1)[1] if "." in low else ""))
|
||
file_storage.save(raw_path)
|
||
if low.endswith(".zip"):
|
||
_extract_zip(raw_path, staging_dir)
|
||
db_file = os.path.join(staging_dir, "users.db")
|
||
remove_quiet(raw_path)
|
||
elif low.endswith(".db"):
|
||
db_file = os.path.join(staging_dir, "users.db")
|
||
os.replace(raw_path, db_file)
|
||
else:
|
||
raise BackupError("仅支持 .zip(平台导出)或 .db(sqlite 库)文件")
|
||
ok, info = validate_backup(db_file)
|
||
if not ok:
|
||
raise BackupError(info)
|
||
info["file_name"] = orig_name
|
||
info["file_size"] = os.path.getsize(db_file)
|
||
return token, info
|
||
except BackupError:
|
||
remove_quiet(staging_dir)
|
||
raise
|
||
except Exception as e:
|
||
remove_quiet(staging_dir)
|
||
raise BackupError(f"暂存上传文件失败: {e}")
|
||
|
||
|
||
# ================== 导入:应用(落待生效任务) ==================
|
||
def apply_restore(token):
|
||
"""校验 token → 快照当前库 → 把暂存库移到 restore_pending。
|
||
|
||
返回 dict {ok, backup_name, message};失败抛 BackupError。token 形如 12 位 hex。
|
||
"""
|
||
if not token or len(token) != 12 or not all(c in "0123456789abcdef" for c in token):
|
||
raise BackupError("无效的导入标识")
|
||
staging_dir = os.path.join(RESTORE_STAGING_DIR, token)
|
||
db_file = os.path.join(staging_dir, "users.db")
|
||
if not os.path.exists(db_file):
|
||
raise BackupError("导入会话已失效,请重新上传预览")
|
||
ok, info = validate_backup(db_file)
|
||
if not ok:
|
||
raise BackupError(info)
|
||
|
||
# 1) 自动备份当前库(安全网,可回滚)
|
||
_ensure_dir(BACKUP_DIR)
|
||
backup_name = f"pre_restore_{_ts()}.db"
|
||
snapshot_db(os.path.join(BACKUP_DIR, backup_name))
|
||
|
||
# 2) 落待生效恢复任务(旧未消费任务被覆盖——上一份已无意义)
|
||
pending_dir = RESTORE_PENDING_DIR
|
||
if os.path.isdir(pending_dir):
|
||
remove_quiet(pending_dir)
|
||
_ensure_dir(pending_dir)
|
||
os.replace(db_file, os.path.join(pending_dir, "users.db"))
|
||
apk_src = os.path.join(staging_dir, "apks")
|
||
if os.path.isdir(apk_src) and os.listdir(apk_src):
|
||
shutil.move(apk_src, os.path.join(pending_dir, "apks"))
|
||
remove_quiet(staging_dir)
|
||
|
||
schema_version = info.get("schema_version", 0)
|
||
_log.info(f"导入已落恢复任务: schema v{schema_version},"
|
||
f"重启 web_server 后生效;当前库已快照 {backup_name}")
|
||
return {
|
||
"ok": True,
|
||
"backup_name": backup_name,
|
||
"message": "恢复任务已生成:当前库已自动备份,重启 web_server 后即应用导入的数据",
|
||
}
|
||
|
||
|
||
# ================== 启动消费(web_server.py init_db 前调用) ==================
|
||
def consume_pending_restore():
|
||
"""把 restore_pending/users.db 换位为当前 users.db(删除旧 wal/shm,合并 apks)。
|
||
|
||
必须在 SQLAlchemy engine 首次打开 users.db 之前执行(web_server 装配时调用)。
|
||
失败不阻塞启动:坏恢复文件会挪到 BACKUP_DIR/restore_failed_*/ 并继续用旧库。
|
||
返回是否真的应用了恢复。
|
||
"""
|
||
pending_db = os.path.join(RESTORE_PENDING_DIR, "users.db")
|
||
if not os.path.exists(pending_db):
|
||
return False
|
||
# 消费前轻量复验(防止意外损坏文件把主库换掉)
|
||
ok, msg = validate_backup(pending_db)
|
||
if not ok:
|
||
fail_dir = os.path.join(BACKUP_DIR, f"restore_failed_{_ts()}")
|
||
_ensure_dir(fail_dir)
|
||
os.replace(pending_db, os.path.join(fail_dir, "users.db"))
|
||
remove_quiet(RESTORE_PENDING_DIR)
|
||
_log.error(f"恢复任务校验失败,已搁置到 {fail_dir},继续使用当前库: {msg}")
|
||
return False
|
||
_ensure_dir(DATA_DIR)
|
||
os.replace(pending_db, DB_FILE)
|
||
for ext in ("-wal", "-shm"):
|
||
p = DB_FILE + ext
|
||
if os.path.exists(p):
|
||
remove_quiet(p)
|
||
# 合并 apks(覆盖同名)
|
||
apk_src = os.path.join(RESTORE_PENDING_DIR, "apks")
|
||
if os.path.isdir(apk_src):
|
||
_ensure_dir(APK_DIR)
|
||
for fname in os.listdir(apk_src):
|
||
src = os.path.join(apk_src, fname)
|
||
dst = os.path.join(APK_DIR, fname)
|
||
if os.path.exists(dst):
|
||
os.remove(dst)
|
||
shutil.move(src, dst)
|
||
remove_quiet(RESTORE_PENDING_DIR)
|
||
schema_version = msg.get("schema_version", 0) if isinstance(msg, dict) else "?"
|
||
_log.info(f"已应用备份恢复: users.db 替换完成(schema v{schema_version}),apks 已合并")
|
||
return True
|