a3216 commited on
Commit
48d810d
·
verified ·
1 Parent(s): 58947b2

sync from GitHub 1c1ef96: feat: implement v6 backup storage layer with content-addressable blobs, snapshot management, and garbage collection

Browse files
app/api/user_backups_v6.py ADDED
@@ -0,0 +1,393 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """v6 备份 API:/u/api/backups/v6/*
2
+
3
+ 端点(12 个):
4
+ - GET /sets 列出当前用户备份点
5
+ - POST /sets 创建待提交备份点
6
+ - GET /sets/{setId} 获取备份点 manifest
7
+ - DELETE /sets/{setId} 删除备份点(快照模型:任意点可删)
8
+ - POST /sets/{setId}/manifest 提交 manifest(原子生效)
9
+ - POST /sets/{setId}/rename 备份点改名
10
+ - GET /sets/{setId}/archive/info ZIP 分卷信息
11
+ - GET /sets/{setId}/archive?volume=N 下载第 N 卷 ZIP(request.download 通道)
12
+ - POST /blobs/upload 上传单个 blob(multipart,小文件)
13
+ - POST /blobs/upload_part 大文件 base64 分片上传
14
+ - POST /blobs/batch_check 卸载重装兜底(path+size+mtime 比对最新备份点)
15
+ - DELETE /legacy-v5 清理当前用户的 v5 遗留数据(被污染,整体废弃)
16
+ """
17
+ from __future__ import annotations
18
+
19
+ import asyncio
20
+ import logging
21
+
22
+ from fastapi import APIRouter, File, Form, Request, UploadFile
23
+ from fastapi.responses import FileResponse, StreamingResponse
24
+ from typing import Optional
25
+
26
+ from ..errors import HttpError
27
+ from ..services import backup_v6_store
28
+ from ._common import CORS_HEADERS, ok_with_cors, read_json_body
29
+ from .user_files import require_user, require_user_for_download
30
+
31
+ logger = logging.getLogger(__name__)
32
+
33
+ router = APIRouter(prefix="/u/api/backups/v6", tags=["user-backups-v6"])
34
+
35
+ # 单 blob 上限:框架单文件硬限制 10MB,但恢复走 ZIP 分卷(每卷预算 8MB),
36
+ # 单 blob 必须给卷内 manifest.json 留余量,与前端 backup-v6.js 的 MAX_BLOB_SIZE 一致
37
+ MAX_BLOB_SIZE = 7 * 1024 * 1024
38
+
39
+ # 后台 Hub 同步任务引用(保存强引用防 task 被 GC 中途回收)
40
+ _sync_task_ref = None
41
+
42
+
43
+ def _err(code: str, message: str, status: int = 400, **extra) -> HttpError:
44
+ return HttpError(message, status=status, code=code, details=extra or None)
45
+
46
+
47
+ def _user_id(request: Request) -> str:
48
+ payload = require_user(request)
49
+ uid = str(payload.get("sub") or "")
50
+ if not uid:
51
+ raise _err("unauthorized", "missing user", status=401)
52
+ return uid
53
+
54
+
55
+ # ===== sets =====
56
+
57
+ @router.get("/sets")
58
+ async def list_sets(request: Request):
59
+ uid = _user_id(request)
60
+ try:
61
+ limit = int(request.query_params.get("limit") or "30")
62
+ offset = int(request.query_params.get("offset") or "0")
63
+ except Exception:
64
+ limit, offset = 30, 0
65
+ # 本地 DB 为空时触发 Hub 同步(HF Space 重建后 DB 丢失)
66
+ if backup_v6_store.count_sets(uid) == 0:
67
+ global _sync_task_ref
68
+ if _sync_task_ref is None or _sync_task_ref.done():
69
+ _sync_task_ref = asyncio.create_task(backup_v6_store._sync_sets_background())
70
+ items = backup_v6_store.list_sets(uid, limit=limit, offset=offset)
71
+ return ok_with_cors({"items": items, "count": len(items), "total": len(items)})
72
+
73
+
74
+ @router.post("/sets")
75
+ async def create_set(request: Request):
76
+ uid = _user_id(request)
77
+ body = await read_json_body(request)
78
+ alias = str(body.get("alias") or "").strip()[:200]
79
+ try:
80
+ result = backup_v6_store.create_pending_set(uid, alias=alias)
81
+ except ValueError as e:
82
+ msg = str(e)
83
+ if msg == "quota_exceeded":
84
+ raise _err("quota_exceeded", f"已达备份点上限({backup_v6_store.MAX_SETS_PER_USER} 份),请先删除旧的备份", status=409)
85
+ raise _err("create_failed", msg, status=500)
86
+ # 小天才 fetch.fetch 对非 200 的 2xx(如 201)会触发 fail 回调,统一返回 200
87
+ return ok_with_cors(result, status=200)
88
+
89
+
90
+ @router.get("/sets/{set_id}")
91
+ async def get_set(set_id: str, request: Request):
92
+ uid = _user_id(request)
93
+ owner = backup_v6_store.get_set_owner(set_id)
94
+ if not owner:
95
+ raise _err("set_not_found", "备份点不存在", status=404)
96
+ if owner != uid:
97
+ raise _err("forbidden", "无权访问他人备份", status=403)
98
+ manifest = backup_v6_store.get_set_manifest(set_id)
99
+ if manifest is None:
100
+ raise _err("manifest_missing", "manifest 缺失", status=500)
101
+ return ok_with_cors({"manifest": manifest})
102
+
103
+
104
+ @router.delete("/sets/{set_id}")
105
+ async def delete_set(set_id: str, request: Request):
106
+ uid = _user_id(request)
107
+ owner = backup_v6_store.get_set_owner(set_id)
108
+ if not owner:
109
+ raise _err("set_not_found", "备份点不存在", status=404)
110
+ if owner != uid:
111
+ raise _err("forbidden", "无权删除他人备份", status=403)
112
+ try:
113
+ result = backup_v6_store.delete_set(set_id, uid)
114
+ except ValueError as e:
115
+ msg = str(e)
116
+ if msg == "set_not_found":
117
+ raise _err("set_not_found", "备份点不存在", status=404)
118
+ if msg == "forbidden":
119
+ raise _err("forbidden", "无权删除他人备份", status=403)
120
+ raise _err("delete_failed", msg, status=500)
121
+ return ok_with_cors({"deleted": True, "setId": set_id, **result})
122
+
123
+
124
+ @router.post("/sets/{set_id}/manifest")
125
+ async def commit_manifest(set_id: str, request: Request):
126
+ uid = _user_id(request)
127
+ body = await read_json_body(request)
128
+ manifest = body.get("manifest")
129
+ if not isinstance(manifest, dict):
130
+ raise _err("invalid_manifest", "manifest 必须是 JSON 对象")
131
+ if str(manifest.get("setId") or "") != set_id:
132
+ manifest["setId"] = set_id
133
+ blobs = manifest.get("blobs")
134
+ if not isinstance(blobs, list) or not blobs:
135
+ raise _err("invalid_manifest_blobs", "manifest.blobs 必须是非空数组")
136
+ keys = [str(b.get("key") or "") for b in blobs if isinstance(b, dict)]
137
+ if not keys:
138
+ raise _err("empty_blobs", "manifest.blobs 不能为空")
139
+ # blob 可用性校验(含 Hub 回源,网络往返放工作线程避免阻塞事件循环)
140
+ try:
141
+ still_missing = await asyncio.to_thread(backup_v6_store.validate_blob_keys, keys)
142
+ except Exception as e:
143
+ raise _err("validate_failed", str(e), status=500)
144
+ if still_missing:
145
+ # 本地与 Hub 都没有:清理孤儿 DB 记录,让客户端重传
146
+ await asyncio.to_thread(backup_v6_store.cleanup_missing_blob_records, still_missing)
147
+ raise _err(
148
+ "blob_missing", "manifest 引用了未上传的 blob",
149
+ missing_keys=still_missing[:20],
150
+ )
151
+ try:
152
+ result = backup_v6_store.commit_set(set_id, manifest, uid)
153
+ except ValueError as e:
154
+ msg = str(e)
155
+ if msg == "set_not_pending":
156
+ raise _err("set_not_pending", "备份点不在待提交状态(已提交或已过期)")
157
+ if msg == "invalid_manifest":
158
+ raise _err("invalid_manifest", "manifest 格式无效")
159
+ if msg == "invalid_manifest_version":
160
+ raise _err("invalid_manifest_version", "manifest.version 必须为 6")
161
+ if msg in ("invalid_manifest_blobs", "empty_blobs"):
162
+ raise _err(msg, "manifest.blobs 必须是非空数组")
163
+ if msg.startswith("blob_missing:"):
164
+ missing = msg.split(":", 1)[1].split(",")
165
+ raise _err("blob_missing", "manifest 引用了未上传的 blob", missing_keys=missing)
166
+ if msg == "quota_exceeded":
167
+ raise _err("quota_exceeded", "已达备份点上限", status=409)
168
+ raise _err("commit_failed", msg, status=500)
169
+ return ok_with_cors(result)
170
+
171
+
172
+ @router.post("/sets/{set_id}/rename")
173
+ async def rename_set(set_id: str, request: Request):
174
+ uid = _user_id(request)
175
+ body = await read_json_body(request)
176
+ alias = str(body.get("alias") or "").strip()[:200]
177
+ if not alias:
178
+ raise _err("invalid_alias", "别名不能为空")
179
+ try:
180
+ result = backup_v6_store.rename_set(set_id, uid, alias)
181
+ except ValueError as e:
182
+ msg = str(e)
183
+ if msg == "set_not_found":
184
+ raise _err("set_not_found", "备份点不存在", status=404)
185
+ if msg == "forbidden":
186
+ raise _err("forbidden", "无权操作他人备份", status=403)
187
+ raise _err("rename_failed", msg, status=500)
188
+ return ok_with_cors(result)
189
+
190
+
191
+ # ===== ZIP 分卷下载 =====
192
+
193
+ @router.get("/sets/{set_id}/archive/info")
194
+ async def archive_info(set_id: str, request: Request):
195
+ uid = _user_id(request)
196
+ owner = backup_v6_store.get_set_owner(set_id)
197
+ if not owner:
198
+ raise _err("set_not_found", "备份点不存在", status=404)
199
+ if owner != uid:
200
+ raise _err("forbidden", "无权访问他人备份", status=403)
201
+ try:
202
+ info = backup_v6_store.archive_info(set_id)
203
+ except ValueError as e:
204
+ if str(e) == "set_not_found":
205
+ raise _err("set_not_found", "备份点不存在", status=404)
206
+ raise _err("archive_info_failed", str(e), status=500)
207
+ return ok_with_cors(info)
208
+
209
+
210
+ @router.get("/sets/{set_id}/archive")
211
+ async def download_archive(set_id: str, request: Request):
212
+ payload = require_user_for_download(request)
213
+ uid = str(payload.get("sub") or "")
214
+ if not uid:
215
+ raise _err("unauthorized", "missing user", status=401)
216
+ owner = backup_v6_store.get_set_owner(set_id)
217
+ if not owner:
218
+ raise _err("set_not_found", "备份点不存在", status=404)
219
+ if owner != uid:
220
+ raise _err("forbidden", "无权下载他人备份", status=403)
221
+ try:
222
+ volume = int(request.query_params.get("volume") or "0")
223
+ except Exception:
224
+ volume = 0
225
+ try:
226
+ zip_path, zip_size = await asyncio.to_thread(
227
+ backup_v6_store.build_volume_zip, set_id, volume
228
+ )
229
+ except ValueError as e:
230
+ msg = str(e)
231
+ if msg == "set_not_found":
232
+ raise _err("set_not_found", "备份点不存在", status=404)
233
+ if msg == "volume_out_of_range":
234
+ raise _err("volume_out_of_range", "分卷号超出范围", status=400)
235
+ raise _err("zip_build_failed", msg, status=500)
236
+
237
+ filename = f"backup_{set_id[:8]}_v{volume}.zip"
238
+ generator = backup_v6_store.zip_iter_chunks(zip_path)
239
+ headers = {
240
+ "Content-Disposition": f'attachment; filename="{filename}"',
241
+ "Content-Length": str(zip_size),
242
+ **CORS_HEADERS,
243
+ }
244
+ return StreamingResponse(generator, media_type="application/zip", headers=headers)
245
+
246
+
247
+ # ===== blobs =====
248
+
249
+ @router.post("/blobs/upload")
250
+ async def upload_blob(
251
+ request: Request,
252
+ file: UploadFile = File(...),
253
+ path: str = Form(""),
254
+ type: str = Form(""),
255
+ ):
256
+ uid = _user_id(request)
257
+ content = await file.read()
258
+ if len(content) > MAX_BLOB_SIZE:
259
+ raise _err("blob_too_large", f"blob 超过 {MAX_BLOB_SIZE} 字节", status=413)
260
+ mime = str(file.content_type or "")
261
+ if not mime:
262
+ mime = backup_v6_store._guess_mime(path or file.filename or "")
263
+ try:
264
+ result = backup_v6_store.store_blob(content, mime)
265
+ except ValueError as e:
266
+ if str(e) == "blob_too_large":
267
+ raise _err("blob_too_large", "blob 过大", status=413)
268
+ raise _err("store_failed", str(e), status=500)
269
+ return ok_with_cors(result)
270
+
271
+
272
+ @router.post("/blobs/upload_part")
273
+ async def upload_blob_part(request: Request):
274
+ """大文件分片上传:{uploadId, seq, dataB64, last} → {key, size, exists}。
275
+
276
+ 流程:POST /blobs/upload_part/start 拿 uploadId → 逐片 POST → last=true 收尾。
277
+ """
278
+ uid = _user_id(request)
279
+ body = await read_json_body(request)
280
+ upload_id = str(body.get("uploadId") or "")
281
+ seq = body.get("seq")
282
+ data_b64 = body.get("dataB64")
283
+ last = body.get("last") is True
284
+ if not upload_id:
285
+ raise _err("invalid_upload_id", "uploadId 不能为空")
286
+ try:
287
+ seq = int(seq)
288
+ except Exception:
289
+ raise _err("invalid_seq", "seq 必须是整数")
290
+ if not isinstance(data_b64, str) or not data_b64:
291
+ raise _err("invalid_data", "dataB64 不能为空")
292
+ try:
293
+ result = backup_v6_store.upload_part(upload_id, uid, seq, data_b64, last)
294
+ except ValueError as e:
295
+ msg = str(e)
296
+ if msg == "upload_not_found":
297
+ raise _err("upload_not_found", "上传会话不存在或已过期", status=404)
298
+ if msg == "upload_expired":
299
+ raise _err("upload_expired", "上传会话已过期,请重新开始")
300
+ if msg == "part_too_large":
301
+ raise _err("part_too_large", f"单分片解码后超过 {backup_v6_store.MAX_PART_SIZE} 字节", status=413)
302
+ if msg == "invalid_base64":
303
+ raise _err("invalid_base64", "dataB64 不是合法 base64")
304
+ raise _err("upload_failed", msg, status=500)
305
+ return ok_with_cors(result)
306
+
307
+
308
+ @router.post("/blobs/upload_part/start")
309
+ async def upload_blob_part_start(request: Request):
310
+ uid = _user_id(request)
311
+ body = await read_json_body(request)
312
+ path = str(body.get("path") or "")
313
+ type_ = str(body.get("type") or "")
314
+ if not path:
315
+ raise _err("invalid_path", "path 不能为空")
316
+ upload_id = backup_v6_store.start_upload(uid, path, type_)
317
+ return ok_with_cors({"uploadId": upload_id, "expiresIn": backup_v6_store.PENDING_SET_TTL_SEC})
318
+
319
+
320
+ @router.post("/blobs/batch_check")
321
+ async def batch_check(request: Request):
322
+ uid = _user_id(request)
323
+ body = await read_json_body(request)
324
+ paths = body.get("paths")
325
+ if not isinstance(paths, list):
326
+ raise _err("invalid_paths", "paths 必须是数组")
327
+ if len(paths) > backup_v6_store.BATCH_CHECK_MAX:
328
+ raise _err("too_many_paths", f"单次最多 {backup_v6_store.BATCH_CHECK_MAX} 个 path")
329
+ cleaned = []
330
+ for p in paths:
331
+ if not isinstance(p, dict):
332
+ continue
333
+ path = str(p.get("path") or "").strip()
334
+ if not path:
335
+ continue
336
+ try:
337
+ size = int(p.get("size") or 0)
338
+ except Exception:
339
+ size = 0
340
+ try:
341
+ mtime = int(p.get("mtime") or 0)
342
+ except Exception:
343
+ mtime = 0
344
+ cleaned.append({"path": path, "size": size, "mtime": mtime})
345
+ result = backup_v6_store.batch_check(uid, cleaned)
346
+ return ok_with_cors(result)
347
+
348
+
349
+ @router.get("/blobs/{key}")
350
+ async def download_blob(key: str, request: Request):
351
+ payload = require_user_for_download(request)
352
+ uid = str(payload.get("sub") or "")
353
+ if not uid:
354
+ raise _err("unauthorized", "missing user", status=401)
355
+ if not _is_sha256_hex(key):
356
+ raise _err("invalid_key", "key 必须是 64 位 sha256 hex", status=400)
357
+ blob_path = backup_v6_store.get_blob_local_path(key)
358
+ if not blob_path or not blob_path.exists():
359
+ raise _err("blob_not_found", "blob 不存在", status=404)
360
+ meta = backup_v6_store.get_blob_meta(key)
361
+ mime = (meta or {}).get("mime") or "application/octet-stream"
362
+ headers = {
363
+ "Content-Disposition": f'attachment; filename="{blob_path.name}"',
364
+ **CORS_HEADERS,
365
+ }
366
+ return FileResponse(str(blob_path), media_type=mime, headers=headers)
367
+
368
+
369
+ # ===== v5 遗留数据清理 =====
370
+
371
+ @router.delete("/legacy-v5")
372
+ async def purge_legacy_v5(request: Request):
373
+ """清理当前用户的 v5 遗留数据(被污染 + 设计缺陷,已整体废弃)。
374
+
375
+ 删除该用户全部 v5 备份集、refs=0 的 v5 blob(本地 + Hub)。
376
+ v6 数据不受任何影响。
377
+ """
378
+ uid = _user_id(request)
379
+ try:
380
+ result = await asyncio.to_thread(backup_v6_store.purge_legacy_v5, uid)
381
+ except Exception as e:
382
+ raise _err("purge_failed", str(e), status=500)
383
+ return ok_with_cors({"ok": True, "purged": result})
384
+
385
+
386
+ def _is_sha256_hex(s: str) -> bool:
387
+ if len(s) != 64:
388
+ return False
389
+ try:
390
+ int(s, 16)
391
+ return True
392
+ except Exception:
393
+ return False
app/database.py CHANGED
@@ -336,6 +336,35 @@ def _create_schema(conn: sqlite3.Connection) -> None:
336
  alias TEXT
337
  );
338
  CREATE INDEX IF NOT EXISTS idx_backup_sets_user ON backup_sets(user_id, created_at DESC);
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
339
  """
340
  )
341
  # 兼容已存在的库:若旧表无 device_id 列则补加(ALTER TABLE 幂等检查)
 
336
  alias TEXT
337
  );
338
  CREATE INDEX IF NOT EXISTS idx_backup_sets_user ON backup_sets(user_id, created_at DESC);
339
+
340
+ -- v6 备份:内容寻址 blob(全局共享,refs 引用计数)
341
+ -- 与 v5 分表:v5 数据已被污染,v6 独立命名空间便于整体清理旧数据
342
+ CREATE TABLE IF NOT EXISTS backup_v6_blobs (
343
+ key TEXT PRIMARY KEY,
344
+ size INTEGER NOT NULL,
345
+ mime TEXT,
346
+ refs INTEGER NOT NULL DEFAULT 0,
347
+ uploaded_at INTEGER NOT NULL
348
+ );
349
+ CREATE INDEX IF NOT EXISTS idx_backup_v6_blobs_refs ON backup_v6_blobs(refs);
350
+
351
+ -- v6 备份:备份点(快照模型:每个 set 的 manifest 都列全量文件清单,
352
+ -- blob 跨 set 共享。删除任意 set 只减 refs,其他 set 完全不受影响)
353
+ CREATE TABLE IF NOT EXISTS backup_v6_sets (
354
+ set_id TEXT PRIMARY KEY,
355
+ user_id TEXT NOT NULL,
356
+ alias TEXT,
357
+ manifest TEXT NOT NULL,
358
+ manifest_sha256 TEXT NOT NULL,
359
+ total_size INTEGER NOT NULL,
360
+ blob_count INTEGER NOT NULL,
361
+ created_at INTEGER NOT NULL,
362
+ is_latest INTEGER NOT NULL DEFAULT 0,
363
+ prev_set_id TEXT,
364
+ mode TEXT,
365
+ device TEXT
366
+ );
367
+ CREATE INDEX IF NOT EXISTS idx_backup_v6_sets_user ON backup_v6_sets(user_id, created_at DESC);
368
  """
369
  )
370
  # 兼容已存在的库:若旧表无 device_id 列则补加(ALTER TABLE 幂等检查)
app/main.py CHANGED
@@ -150,6 +150,10 @@ async def lifespan(app: FastAPI):
150
  # v5 备份:从 Hub 同步 manifest 元数据回本地 DB(HF Space 重建后本地 DB 丢失,
151
  # 但 Hub 上 backup_sets_v5/<user_id>/<set_id>.json 仍存在,需恢复否则管理面板和列表为空)
152
  asyncio.create_task(backup_v5_store._sync_manifests_background())
 
 
 
 
153
  # 启动用户文件 Hub 定时同步任务
154
  s = get_settings()
155
  if is_hub_enabled():
@@ -176,6 +180,9 @@ async def lifespan(app: FastAPI):
176
  # 停止 v5 备份 GC 任务
177
  from .services import backup_v5_store
178
  await backup_v5_store.stop_gc_task()
 
 
 
179
  if sync_task is not None:
180
  await user_files_sync.stop_sync_task()
181
  await close_http_client()
@@ -266,6 +273,7 @@ def create_app() -> FastAPI:
266
  user_files,
267
  user_backups,
268
  user_backups_v5,
 
269
  webhooks,
270
  xtc,
271
  )
@@ -289,6 +297,7 @@ def create_app() -> FastAPI:
289
  app.include_router(user_files.router)
290
  app.include_router(user_backups.router)
291
  app.include_router(user_backups_v5.router)
 
292
  app.include_router(music.router)
293
  app.include_router(music_records.router)
294
  app.include_router(logs.router)
 
150
  # v5 备份:从 Hub 同步 manifest 元数据回本地 DB(HF Space 重建后本地 DB 丢失,
151
  # 但 Hub 上 backup_sets_v5/<user_id>/<set_id>.json 仍存在,需恢复否则管理面板和列表为空)
152
  asyncio.create_task(backup_v5_store._sync_manifests_background())
153
+ # v6 备份:GC 任务(refs=0 且超宽限期的 blob)+ Hub 元数据同步(Space 重建恢复)
154
+ from .services import backup_v6_store
155
+ backup_v6_store.start_gc_task()
156
+ asyncio.create_task(backup_v6_store._sync_sets_background())
157
  # 启动用户文件 Hub 定时同步任务
158
  s = get_settings()
159
  if is_hub_enabled():
 
180
  # 停止 v5 备份 GC 任务
181
  from .services import backup_v5_store
182
  await backup_v5_store.stop_gc_task()
183
+ # 停止 v6 备份 GC 任务
184
+ from .services import backup_v6_store
185
+ await backup_v6_store.stop_gc_task()
186
  if sync_task is not None:
187
  await user_files_sync.stop_sync_task()
188
  await close_http_client()
 
273
  user_files,
274
  user_backups,
275
  user_backups_v5,
276
+ user_backups_v6,
277
  webhooks,
278
  xtc,
279
  )
 
297
  app.include_router(user_files.router)
298
  app.include_router(user_backups.router)
299
  app.include_router(user_backups_v5.router)
300
+ app.include_router(user_backups_v6.router)
301
  app.include_router(music.router)
302
  app.include_router(music_records.router)
303
  app.include_router(logs.router)
app/services/backup_v6_store.py ADDED
@@ -0,0 +1,1098 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """v6 备份存储层:内容寻址 blob + 快照备份点 + ZIP 分卷 + GC + v5 数据清理。
2
+
3
+ 与 v5 的核心差异(针对 v5 实际运行中暴露的顽固问题重新设计):
4
+ 1. 快照模型:每个备份点的 manifest 都列全量文件清单,blob 跨备份点共享。
5
+ 删除任意备份点(包括最初那份)只减引用计数,其他备份点完全不受影响,
6
+ 任意备份点可独立恢复。v5 的 prevSetId 增量链导致的"删中间点链断裂"
7
+ 问题从模型上消除。
8
+ 2. 分片上传通道:≥512KB 的大文件可走 base64 分片(fetch.fetch JSON),
9
+ 规避 request.upload multipart 对大文件回调不稳定的平台缺陷。
10
+ 3. ZIP 分卷打包:单卷 ≤8MB(框架单文件 10MB 硬限制留余量),大数据量
11
+ 备份不再因"超 10MB 无法 ZIP"而被迫走会耗尽下载管理器的逐 blob 通道。
12
+ ZIP 内保持扁平结构 blob_<idx>.dat(小天才 zip.decompress 对子目录
13
+ ZIP 解压不可靠,v5 实测结论)。
14
+ 4. blob 双重校验:commit 时校验 DB 记录 + 本地文件,本地丢失自动从 Hub
15
+ 拉回(HF Space 重建场景)。
16
+ 5. 独立命名空间 backup_v6_blobs / backup_v6_sets:与被污染的 v5 数据
17
+ 完全隔离,v5 数据可整体清理(purge_legacy_v5)。
18
+
19
+ 存储布局:
20
+ - 本地:/data/hf_bucket/backup_v6_blobs/<key[:2]>/<key>
21
+ /data/hf_bucket/backup_v6_sets/<user_id>/<set_id>.json
22
+ - Hub: 同名空间(冷备,HF Space 重建后恢复)
23
+ - DB: backup_v6_blobs / backup_v6_sets 表
24
+ """
25
+ from __future__ import annotations
26
+
27
+ import asyncio
28
+ import base64
29
+ import hashlib
30
+ import json
31
+ import logging
32
+ import time
33
+ import zipfile
34
+ from pathlib import Path
35
+ from typing import Optional
36
+
37
+ from ..database import get_conn
38
+ from ..hf_storage import (
39
+ delete_from_hub,
40
+ download_bytes,
41
+ is_hub_enabled,
42
+ list_hub_keys,
43
+ push_to_hub,
44
+ upload_bytes,
45
+ )
46
+
47
+ logger = logging.getLogger(__name__)
48
+
49
+ # ===== 常量 =====
50
+
51
+ NS_BLOBS = "backup_v6_blobs"
52
+ NS_SETS = "backup_v6_sets"
53
+
54
+ MAX_BLOB_SIZE = 10 * 1024 * 1024 # 单 blob ≤10MB(框架硬限制)
55
+ MAX_PART_SIZE = 1024 * 1024 # 单分片解码后 ≤1MB
56
+ VOLUME_BUDGET = 8 * 1024 * 1024 # ZIP 单卷内容预算 ≤8MB(10MB 限制留余量)
57
+ MAX_SETS_PER_USER = 8 # 单用户保留最近 8 个备份点
58
+ PENDING_SET_TTL_SEC = 600 # 待提交 set / 分片上传 TTL:10 分钟
59
+ GC_INTERVAL_SEC = 3600 # GC 周期:1 小时
60
+ GC_BLOB_GRACE_SEC = 24 * 3600 # refs=0 的 blob 保留 24h 再删(防误删进行中的上传)
61
+ BATCH_CHECK_MAX = 300 # batch_check 单次最大 path 数
62
+
63
+ # 模块级内存状态(多 worker 场景各自独立,可接受:失败重传即可)
64
+ _pending_sets: "dict[str, dict]" = {} # set_id -> {user_id, alias, created_at, prev_set_id}
65
+ _pending_uploads: "dict[str, dict]" = {} # uploadId -> {user_id, parts, created_at, path, type}
66
+
67
+
68
+ def _now() -> int:
69
+ return int(time.time())
70
+
71
+
72
+ def _blob_local_path(key: str) -> Path:
73
+ from ..hf_storage import _local_path
74
+ return _local_path(NS_BLOBS, f"{key[:2]}/{key}")
75
+
76
+
77
+ def _manifest_local_path(user_id: str, set_id: str) -> Path:
78
+ from ..hf_storage import _local_path
79
+ return _local_path(NS_SETS, f"{user_id}/{set_id}.json")
80
+
81
+
82
+ def _gen_id() -> str:
83
+ import uuid
84
+ return uuid.uuid4().hex
85
+
86
+
87
+ def _cleanup_expired() -> None:
88
+ """清理过期的待提交 set 与分片上传缓冲。"""
89
+ now = _now()
90
+ expired = [sid for sid, p in _pending_sets.items() if now - p["created_at"] > PENDING_SET_TTL_SEC]
91
+ for sid in expired:
92
+ _pending_sets.pop(sid, None)
93
+ expired2 = [uid for uid, p in _pending_uploads.items() if now - p["created_at"] > PENDING_SET_TTL_SEC]
94
+ for uid in expired2:
95
+ _pending_uploads.pop(uid, None)
96
+
97
+
98
+ # ===== Blob 存储(内容寻址)=====
99
+
100
+ def store_blob(content: bytes, mime: str) -> dict:
101
+ """存 blob:key = sha256(content)。已存在则去重。
102
+
103
+ 返回 {key, size, exists}。不增 refs(refs 在 commit_set 时统一累加)。
104
+ """
105
+ if len(content) > MAX_BLOB_SIZE:
106
+ raise ValueError("blob_too_large")
107
+ key = hashlib.sha256(content).hexdigest()
108
+ size = len(content)
109
+ now = _now()
110
+ p = _blob_local_path(key)
111
+
112
+ with get_conn() as conn:
113
+ row = conn.execute(
114
+ "SELECT key FROM backup_v6_blobs WHERE key = ?", (key,)
115
+ ).fetchone()
116
+ if row:
117
+ # DB 有记录但本地文件丢失(HF Space 重建后常见):补写本地,
118
+ # 否则 commit 校验会报 blob_missing(v5 踩过的坑)
119
+ if not p.exists():
120
+ try:
121
+ p.parent.mkdir(parents=True, exist_ok=True)
122
+ p.write_bytes(content)
123
+ if is_hub_enabled():
124
+ asyncio.ensure_future(_push_blob_to_hub(key))
125
+ except Exception as e:
126
+ logger.warning("[v6_blob] restore local file failed key=%s: %s", key[:8], e)
127
+ return {"key": key, "size": size, "exists": True}
128
+
129
+ p.parent.mkdir(parents=True, exist_ok=True)
130
+ p.write_bytes(content)
131
+ try:
132
+ with get_conn() as conn:
133
+ conn.execute(
134
+ "INSERT OR IGNORE INTO backup_v6_blobs(key, size, mime, refs, uploaded_at) "
135
+ "VALUES(?, ?, ?, 0, ?)",
136
+ (key, size, mime, now),
137
+ )
138
+ except Exception as e:
139
+ logger.warning("[v6_blob] insert db failed key=%s: %s", key[:8], e)
140
+ if is_hub_enabled():
141
+ asyncio.ensure_future(_push_blob_to_hub(key))
142
+ return {"key": key, "size": size, "exists": False}
143
+
144
+
145
+ async def _push_blob_to_hub(key: str) -> None:
146
+ p = _blob_local_path(key)
147
+ if not p.exists():
148
+ return
149
+ try:
150
+ sub_key = f"{key[:2]}/{key}"
151
+ upload_bytes(NS_BLOBS, sub_key, p.read_bytes())
152
+ await push_to_hub(NS_BLOBS, sub_key)
153
+ except Exception as e:
154
+ logger.warning("[v6_blob] push to hub failed key=%s: %s", key[:8], e)
155
+
156
+
157
+ def get_blob_local_path(key: str) -> Optional[Path]:
158
+ """返回 blob 本地路径;本地缺失时从 Hub 拉回;都没有返回 None。"""
159
+ p = _blob_local_path(key)
160
+ if p.exists() and p.is_file():
161
+ return p
162
+ if is_hub_enabled():
163
+ data = download_bytes(NS_BLOBS, f"{key[:2]}/{key}")
164
+ if data is not None:
165
+ try:
166
+ p.parent.mkdir(parents=True, exist_ok=True)
167
+ p.write_bytes(data)
168
+ return p
169
+ except Exception as e:
170
+ logger.warning("[v6_blob] fetch from hub failed key=%s: %s", key[:8], e)
171
+ return None
172
+
173
+
174
+ def get_blob_meta(key: str) -> Optional[dict]:
175
+ with get_conn() as conn:
176
+ row = conn.execute(
177
+ "SELECT key, size, mime, refs, uploaded_at FROM backup_v6_blobs WHERE key = ?",
178
+ (key,),
179
+ ).fetchone()
180
+ if not row:
181
+ return None
182
+ d = dict(row)
183
+ d["uploadedAt"] = d.pop("uploaded_at")
184
+ return d
185
+
186
+
187
+ # ===== 分片上传(大文件通道)=====
188
+
189
+ def start_upload(user_id: str, path: str, type_: str) -> str:
190
+ """创建分片上传会话,返回 uploadId。"""
191
+ _cleanup_expired()
192
+ upload_id = _gen_id()
193
+ _pending_uploads[upload_id] = {
194
+ "user_id": user_id,
195
+ "parts": {},
196
+ "created_at": _now(),
197
+ "path": str(path or "")[:256],
198
+ "type": str(type_ or "")[:32],
199
+ }
200
+ return upload_id
201
+
202
+
203
+ def upload_part(upload_id: str, user_id: str, seq: int, data_b64: str, last: bool) -> dict:
204
+ """接收一个分片(base64)。last=true 时拼接所有分片存为 blob。
205
+
206
+ 返回 {key, size, exists}(非 last 时返回 {accepted: true})。
207
+ """
208
+ up = _pending_uploads.get(upload_id)
209
+ if not up or up["user_id"] != user_id:
210
+ raise ValueError("upload_not_found")
211
+ if _now() - up["created_at"] > PENDING_SET_TTL_SEC:
212
+ _pending_uploads.pop(upload_id, None)
213
+ raise ValueError("upload_expired")
214
+
215
+ try:
216
+ raw = base64.b64decode(str(data_b64 or ""), validate=False)
217
+ except Exception:
218
+ raise ValueError("invalid_base64")
219
+ if len(raw) > MAX_PART_SIZE:
220
+ raise ValueError("part_too_large")
221
+ up["parts"][int(seq)] = raw
222
+
223
+ if not last:
224
+ return {"accepted": True}
225
+
226
+ # 拼接所有分片(按 seq 排序)
227
+ content = b"".join(up["parts"][s] for s in sorted(up["parts"].keys()))
228
+ mime = _guess_mime(up["path"])
229
+ result = store_blob(content, mime)
230
+ _pending_uploads.pop(upload_id, None)
231
+ result["size"] = len(content)
232
+ return result
233
+
234
+
235
+ def _guess_mime(path: str) -> str:
236
+ ext = (path or "").rsplit(".", 1)[-1].lower() if "." in (path or "") else ""
237
+ return {
238
+ "json": "application/json",
239
+ "txt": "text/plain",
240
+ "md": "text/markdown",
241
+ "png": "image/png",
242
+ "jpg": "image/jpeg",
243
+ "jpeg": "image/jpeg",
244
+ "gif": "image/gif",
245
+ "webp": "image/webp",
246
+ "mp3": "audio/mpeg",
247
+ "wav": "audio/wav",
248
+ }.get(ext, "application/octet-stream")
249
+
250
+
251
+ # ===== 备份点(快照)管理 =====
252
+
253
+ def count_sets(user_id: str) -> int:
254
+ with get_conn() as conn:
255
+ row = conn.execute(
256
+ "SELECT COUNT(*) AS c FROM backup_v6_sets WHERE user_id = ?", (user_id,)
257
+ ).fetchone()
258
+ return int(row["c"]) if row else 0
259
+
260
+
261
+ def create_pending_set(user_id: str, alias: Optional[str] = None) -> dict:
262
+ """创建待提交备份点,返回 {setId, prevSetId, serverLatest, expiresIn}。"""
263
+ _cleanup_expired()
264
+ with get_conn() as conn:
265
+ cnt = conn.execute(
266
+ "SELECT COUNT(*) AS c FROM backup_v6_sets WHERE user_id = ?", (user_id,)
267
+ ).fetchone()
268
+ if cnt and int(cnt["c"]) >= MAX_SETS_PER_USER:
269
+ raise ValueError("quota_exceeded")
270
+ row = conn.execute(
271
+ "SELECT set_id, created_at, alias, device FROM backup_v6_sets "
272
+ "WHERE user_id = ? AND is_latest = 1 LIMIT 1",
273
+ (user_id,),
274
+ ).fetchone()
275
+ prev_set_id = str(row["set_id"]) if row else ""
276
+ server_latest = None
277
+ if row:
278
+ server_latest = {
279
+ "setId": prev_set_id,
280
+ "exportedAt": int(row["created_at"]),
281
+ "alias": row["alias"] or "",
282
+ "device": row["device"] or "",
283
+ }
284
+ set_id = _gen_id()
285
+ _pending_sets[set_id] = {
286
+ "user_id": user_id,
287
+ "alias": alias or "",
288
+ "created_at": _now(),
289
+ "prev_set_id": prev_set_id,
290
+ }
291
+ return {
292
+ "setId": set_id,
293
+ "prevSetId": prev_set_id,
294
+ "serverLatest": server_latest,
295
+ "expiresIn": PENDING_SET_TTL_SEC,
296
+ }
297
+
298
+
299
+ def get_set_owner(set_id: str) -> Optional[str]:
300
+ with get_conn() as conn:
301
+ row = conn.execute(
302
+ "SELECT user_id FROM backup_v6_sets WHERE set_id = ?", (set_id,)
303
+ ).fetchone()
304
+ return str(row["user_id"]) if row else None
305
+
306
+
307
+ def get_set_manifest(set_id: str) -> Optional[dict]:
308
+ with get_conn() as conn:
309
+ row = conn.execute(
310
+ "SELECT manifest FROM backup_v6_sets WHERE set_id = ?", (set_id,)
311
+ ).fetchone()
312
+ if not row:
313
+ return None
314
+ try:
315
+ return json.loads(row["manifest"])
316
+ except Exception as e:
317
+ logger.warning("[v6_set] parse manifest failed set=%s: %s", set_id[:8], e)
318
+ return None
319
+
320
+
321
+ def list_sets(user_id: str, limit: int = 20, offset: int = 0) -> list[dict]:
322
+ with get_conn() as conn:
323
+ rows = conn.execute(
324
+ "SELECT set_id, manifest, total_size, blob_count, created_at, "
325
+ "is_latest, prev_set_id, alias, mode, device "
326
+ "FROM backup_v6_sets WHERE user_id = ? "
327
+ "ORDER BY created_at DESC LIMIT ? OFFSET ?",
328
+ (user_id, int(limit), int(offset)),
329
+ ).fetchall()
330
+ out = []
331
+ for r in rows:
332
+ try:
333
+ m = json.loads(r["manifest"])
334
+ except Exception:
335
+ m = {}
336
+ stats = m.get("stats", {}) if isinstance(m, dict) else {}
337
+ out.append({
338
+ "setId": r["set_id"],
339
+ "exportedAt": (m.get("exportedAt") if isinstance(m, dict) else None) or int(r["created_at"]),
340
+ "totalSize": int(r["total_size"]),
341
+ "blobCount": int(r["blob_count"]),
342
+ "isLatest": bool(r["is_latest"]),
343
+ "alias": r["alias"] or "",
344
+ "prevSetId": r["prev_set_id"] or "",
345
+ "mode": r["mode"] or "",
346
+ "device": r["device"] or "",
347
+ "stats": stats,
348
+ "volumes": _volume_count(int(r["total_size"])),
349
+ })
350
+ return out
351
+
352
+
353
+ def rename_set(set_id: str, user_id: str, alias: str) -> dict:
354
+ with get_conn() as conn:
355
+ row = conn.execute(
356
+ "SELECT user_id FROM backup_v6_sets WHERE set_id = ?", (set_id,)
357
+ ).fetchone()
358
+ if not row:
359
+ raise ValueError("set_not_found")
360
+ if str(row["user_id"]) != user_id:
361
+ raise ValueError("forbidden")
362
+ conn.execute(
363
+ "UPDATE backup_v6_sets SET alias = ? WHERE set_id = ?",
364
+ (alias[:200], set_id),
365
+ )
366
+ # 同步 Hub 上的 manifest(alias 也写在 manifest 里)
367
+ manifest = get_set_manifest(set_id)
368
+ if manifest is not None:
369
+ manifest["alias"] = alias[:200]
370
+ _write_manifest_to_local(user_id, set_id, manifest)
371
+ if is_hub_enabled():
372
+ asyncio.ensure_future(_push_manifest_to_hub(user_id, set_id))
373
+ return {"setId": set_id, "alias": alias[:200]}
374
+
375
+
376
+ def validate_blob_keys(keys: list) -> list:
377
+ """校验 blob 可用性:DB 记录存在 + 本地文件存在(缺失时从 Hub 拉回)。
378
+
379
+ 返回仍缺失的 key 列表(本地与 Hub 都没有)。
380
+ ★ 含 Hub 网络往返(HF Space 重建后回源),调用方应放 asyncio.to_thread。
381
+ """
382
+ if not keys:
383
+ return []
384
+ uniq = list(dict.fromkeys(k for k in keys if k))
385
+ placeholders = ",".join("?" * len(uniq))
386
+ with get_conn() as conn:
387
+ rows = conn.execute(
388
+ f"SELECT key FROM backup_v6_blobs WHERE key IN ({placeholders})", uniq
389
+ ).fetchall()
390
+ existing = {r["key"] for r in rows}
391
+ missing_db = [k for k in uniq if k not in existing]
392
+ # 本地文件缺失 → 从 Hub 拉回(HF Space 重建场景)
393
+ still_missing = []
394
+ for k in uniq:
395
+ if k in missing_db:
396
+ still_missing.append(k)
397
+ continue
398
+ p = get_blob_local_path(k)
399
+ if not (p and p.exists()):
400
+ still_missing.append(k)
401
+ return still_missing
402
+
403
+
404
+ def cleanup_missing_blob_records(keys: list) -> None:
405
+ """删除本地与 Hub 都缺失的 blob 的 DB 记录(孤儿记录清理)。"""
406
+ for k in keys:
407
+ try:
408
+ with get_conn() as conn:
409
+ conn.execute("DELETE FROM backup_v6_blobs WHERE key = ?", (k,))
410
+ except Exception as e:
411
+ logger.warning("[v6_commit] cleanup orphan record failed key=%s: %s", k[:8], e)
412
+
413
+
414
+ def commit_set(set_id: str, manifest: dict, user_id: str) -> dict:
415
+ """提交 manifest:校验(调用方先 validate_blob_keys)→ 事务写表 → refs 累加
416
+ → 配额淘汰最旧 → 推 Hub。
417
+
418
+ ★ 主写序列(UPDATE is_latest → INSERT set → UPDATE refs)包在
419
+ BEGIN IMMEDIATE 事务里:isolation_level=None 下每语句即时提交,
420
+ 拆开写会在崩溃窗口内造成 refs 少计 → GC 误删在用 blob。
421
+ """
422
+ pending = _pending_sets.get(set_id)
423
+ if not pending or pending["user_id"] != user_id:
424
+ raise ValueError("set_not_pending")
425
+ if not isinstance(manifest, dict):
426
+ raise ValueError("invalid_manifest")
427
+ if int(manifest.get("version") or 0) != 6:
428
+ raise ValueError("invalid_manifest_version")
429
+ blobs = manifest.get("blobs")
430
+ if not isinstance(blobs, list) or not blobs:
431
+ raise ValueError("invalid_manifest_blobs")
432
+
433
+ keys = [str(b.get("key") or "") for b in blobs if isinstance(b, dict)]
434
+ if not keys:
435
+ raise ValueError("empty_blobs")
436
+
437
+ # 校验:DB 记录存在(Hub 回源已由调用方 validate_blob_keys 完成)
438
+ placeholders = ",".join("?" * len(keys))
439
+ with get_conn() as conn:
440
+ rows = conn.execute(
441
+ f"SELECT key FROM backup_v6_blobs WHERE key IN ({placeholders})", keys
442
+ ).fetchall()
443
+ existing_keys = {r["key"] for r in rows}
444
+ missing = [k for k in dict.fromkeys(keys) if k not in existing_keys]
445
+ if missing:
446
+ raise ValueError("blob_missing:" + ",".join(missing[:20]))
447
+
448
+ # checksum
449
+ manifest_for_hash = {k: v for k, v in manifest.items() if k != "checksum"}
450
+ checksum = hashlib.sha256(
451
+ json.dumps(manifest_for_hash, sort_keys=True, ensure_ascii=False).encode("utf-8")
452
+ ).hexdigest()
453
+ manifest["checksum"] = checksum
454
+
455
+ total_size = sum(int(b.get("size") or 0) for b in blobs if isinstance(b, dict))
456
+ blob_count = len(blobs)
457
+ now = _now()
458
+ prev_set_id = manifest.get("prevSetId") or pending.get("prev_set_id") or ""
459
+ alias = (manifest.get("alias") or pending.get("alias") or "")[:200]
460
+ mode = str(manifest.get("mode") or "")[:32]
461
+ device = str(manifest.get("device") or "")[:128]
462
+
463
+ # 配额:>=上限 先删最旧(快照模型:删除不影响其他备份点)
464
+ with get_conn() as conn:
465
+ cnt_row = conn.execute(
466
+ "SELECT COUNT(*) AS c FROM backup_v6_sets WHERE user_id = ?", (user_id,)
467
+ ).fetchone()
468
+ cnt = int(cnt_row["c"]) if cnt_row else 0
469
+ if cnt >= MAX_SETS_PER_USER:
470
+ with get_conn() as conn:
471
+ oldest = conn.execute(
472
+ "SELECT set_id FROM backup_v6_sets WHERE user_id = ? ORDER BY created_at ASC LIMIT 1",
473
+ (user_id,),
474
+ ).fetchone()
475
+ if oldest:
476
+ delete_set(str(oldest["set_id"]), user_id)
477
+
478
+ uniq_keys = list(dict.fromkeys(keys))
479
+ # 事务:is_latest 切换 + INSERT + refs 累加 原子完成
480
+ with get_conn() as conn:
481
+ try:
482
+ conn.execute("BEGIN IMMEDIATE")
483
+ conn.execute(
484
+ "UPDATE backup_v6_sets SET is_latest = 0 WHERE user_id = ? AND is_latest = 1",
485
+ (user_id,),
486
+ )
487
+ conn.execute(
488
+ "INSERT INTO backup_v6_sets(set_id, user_id, alias, manifest, manifest_sha256, "
489
+ "total_size, blob_count, created_at, is_latest, prev_set_id, mode, device) "
490
+ "VALUES(?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?)",
491
+ (
492
+ set_id, user_id, alias, json.dumps(manifest, ensure_ascii=False), checksum,
493
+ total_size, blob_count, now, prev_set_id, mode, device,
494
+ ),
495
+ )
496
+ for k in uniq_keys:
497
+ conn.execute(
498
+ "UPDATE backup_v6_blobs SET refs = refs + 1 WHERE key = ?", (k,)
499
+ )
500
+ conn.execute("COMMIT")
501
+ except Exception:
502
+ try:
503
+ conn.execute("ROLLBACK")
504
+ except Exception:
505
+ pass
506
+ raise
507
+
508
+ _pending_sets.pop(set_id, None)
509
+
510
+ _write_manifest_to_local(user_id, set_id, manifest)
511
+ if is_hub_enabled():
512
+ asyncio.ensure_future(_push_manifest_to_hub(user_id, set_id))
513
+
514
+ return {
515
+ "setId": set_id,
516
+ "manifestSha256": checksum,
517
+ "totalSize": total_size,
518
+ "blobCount": blob_count,
519
+ "isLatest": True,
520
+ }
521
+
522
+
523
+ def delete_set(set_id: str, user_id: str) -> dict:
524
+ """删除备份点:减 refs;refs=0 的 blob 交给 GC(宽限期后清理)。
525
+
526
+ 快照模型下删除任意备份点(含最初那份)都不影响其他备份点。
527
+ """
528
+ with get_conn() as conn:
529
+ row = conn.execute(
530
+ "SELECT user_id, manifest FROM backup_v6_sets WHERE set_id = ?", (set_id,)
531
+ ).fetchone()
532
+ if not row:
533
+ raise ValueError("set_not_found")
534
+ if str(row["user_id"]) != user_id:
535
+ raise ValueError("forbidden")
536
+ try:
537
+ manifest = json.loads(row["manifest"])
538
+ keys = [str(b.get("key") or "") for b in manifest.get("blobs", []) if isinstance(b, dict)]
539
+ except Exception:
540
+ keys = []
541
+
542
+ freed_blobs = 0
543
+ freed_bytes = 0
544
+ with get_conn() as conn:
545
+ for k in dict.fromkeys(keys):
546
+ conn.execute(
547
+ "UPDATE backup_v6_blobs SET refs = MAX(refs - 1, 0) WHERE key = ?", (k,)
548
+ )
549
+ r = conn.execute(
550
+ "SELECT refs, size FROM backup_v6_blobs WHERE key = ?", (k,)
551
+ ).fetchone()
552
+ if r and int(r["refs"]) <= 0:
553
+ freed_blobs += 1
554
+ freed_bytes += int(r["size"])
555
+ conn.execute("DELETE FROM backup_v6_sets WHERE set_id = ?", (set_id,))
556
+
557
+ # 若删的是 latest,把最新的升级为 latest(保持"最新"语义,batch_check 依赖它)
558
+ with get_conn() as conn:
559
+ latest = conn.execute(
560
+ "SELECT set_id FROM backup_v6_sets WHERE user_id = ? AND is_latest = 1 LIMIT 1",
561
+ (user_id,),
562
+ ).fetchone()
563
+ if not latest:
564
+ newest = conn.execute(
565
+ "SELECT set_id FROM backup_v6_sets WHERE user_id = ? ORDER BY created_at DESC LIMIT 1",
566
+ (user_id,),
567
+ ).fetchone()
568
+ if newest:
569
+ conn.execute(
570
+ "UPDATE backup_v6_sets SET is_latest = 1 WHERE set_id = ?",
571
+ (str(newest["set_id"]),),
572
+ )
573
+
574
+ # 删本地 + Hub 的 manifest
575
+ try:
576
+ mp = _manifest_local_path(user_id, set_id)
577
+ if mp.exists():
578
+ mp.unlink(missing_ok=True)
579
+ if is_hub_enabled():
580
+ delete_from_hub(NS_SETS, f"{user_id}/{set_id}.json")
581
+ except Exception as e:
582
+ logger.warning("[v6_set] delete manifest failed set=%s: %s", set_id[:8], e)
583
+
584
+ return {"freedBlobs": freed_blobs, "freedBytes": freed_bytes}
585
+
586
+
587
+ def _write_manifest_to_local(user_id: str, set_id: str, manifest: dict) -> None:
588
+ p = _manifest_local_path(user_id, set_id)
589
+ p.parent.mkdir(parents=True, exist_ok=True)
590
+ p.write_text(json.dumps(manifest, ensure_ascii=False), encoding="utf-8")
591
+
592
+
593
+ async def _push_manifest_to_hub(user_id: str, set_id: str) -> None:
594
+ p = _manifest_local_path(user_id, set_id)
595
+ if not p.exists():
596
+ return
597
+ try:
598
+ sub_key = f"{user_id}/{set_id}.json"
599
+ upload_bytes(NS_SETS, sub_key, p.read_bytes())
600
+ await push_to_hub(NS_SETS, sub_key)
601
+ except Exception as e:
602
+ logger.warning("[v6_set] push manifest to hub failed set=%s: %s", set_id[:8], e)
603
+
604
+
605
+ def batch_check(user_id: str, paths: list[dict]) -> dict:
606
+ """卸载重装兜底:与该用户最新备份点的 manifest 按 path+size+mtime 比对。
607
+
608
+ 返回 {hits: [{path, key}], misses: [{path}]}。
609
+ 注意:比对基准是服务端最新备份点(不需要客户端传 prevSetId)。
610
+ """
611
+ with get_conn() as conn:
612
+ row = conn.execute(
613
+ "SELECT manifest FROM backup_v6_sets WHERE user_id = ? AND is_latest = 1 LIMIT 1",
614
+ (user_id,),
615
+ ).fetchone()
616
+ if not row:
617
+ return {"hits": [], "misses": [{"path": str(p.get("path") or "")} for p in paths]}
618
+ try:
619
+ manifest = json.loads(row["manifest"])
620
+ except Exception:
621
+ manifest = {}
622
+
623
+ prev_map = {}
624
+ for b in manifest.get("blobs", []):
625
+ if isinstance(b, dict):
626
+ prev_map[str(b.get("path") or "")] = {
627
+ "key": str(b.get("key") or ""),
628
+ "size": int(b.get("size") or 0),
629
+ "mtime": int(b.get("mtime") or 0),
630
+ }
631
+
632
+ hits = []
633
+ misses = []
634
+ for p in paths:
635
+ path = str(p.get("path") or "")
636
+ size = int(p.get("size") or 0)
637
+ mtime = int(p.get("mtime") or 0)
638
+ prev = prev_map.get(path)
639
+ if prev and prev["key"] and prev["size"] == size and prev["mtime"] == mtime:
640
+ hits.append({"path": path, "key": prev["key"], "size": size, "mtime": mtime})
641
+ else:
642
+ misses.append({"path": path})
643
+ return {"hits": hits, "misses": misses}
644
+
645
+
646
+ # ===== ZIP 分卷打包 =====
647
+
648
+ def _volume_count(total_size: int) -> int:
649
+ """按 8MB 预算估算卷数(供列表展示)。"""
650
+ if total_size <= 0:
651
+ return 1
652
+ return max(1, (total_size + VOLUME_BUDGET - 1) // VOLUME_BUDGET)
653
+
654
+
655
+ def archive_info(set_id: str) -> dict:
656
+ """返回备份点的分卷信息:{totalSize, volumes, volumeSize, entries}。
657
+
658
+ entries: 每卷的 {volume, startIdx, endIdx(不含), blobCount, contentSize}
659
+ """
660
+ manifest = get_set_manifest(set_id)
661
+ if not manifest:
662
+ raise ValueError("set_not_found")
663
+ # 保持原始数组顺序(分卷下标与 build_volume_zip 的 blob_<i>.dat 对齐),
664
+ # 仅对求和做类型防御(commit 时已校验过全部元素,正常数据无影响)
665
+ blobs = manifest.get("blobs", []) or []
666
+ total_size = sum(
667
+ int(b.get("size") or 0) if isinstance(b, dict) else 0 for b in blobs
668
+ )
669
+
670
+ volumes = []
671
+ cur = None
672
+ for i, b in enumerate(blobs):
673
+ size = int(b.get("size") or 0) if isinstance(b, dict) else 0
674
+ if cur is None or cur["contentSize"] + size > VOLUME_BUDGET:
675
+ if cur is not None:
676
+ volumes.append(cur)
677
+ cur = {"volume": len(volumes), "startIdx": i, "endIdx": i + 1,
678
+ "blobCount": 1, "contentSize": size}
679
+ else:
680
+ cur["endIdx"] = i + 1
681
+ cur["blobCount"] += 1
682
+ cur["contentSize"] += size
683
+ if cur is not None:
684
+ volumes.append(cur)
685
+ if not volumes:
686
+ volumes = [{"volume": 0, "startIdx": 0, "endIdx": 0, "blobCount": 0, "contentSize": 0}]
687
+
688
+ return {
689
+ "setId": set_id,
690
+ "totalSize": total_size,
691
+ "blobCount": len(blobs),
692
+ "volumeBudget": VOLUME_BUDGET,
693
+ "volumeSize": VOLUME_BUDGET,
694
+ "volumes": len(volumes),
695
+ "entries": volumes,
696
+ }
697
+
698
+
699
+ def _zip_tmp_root() -> Path:
700
+ """ZIP 打包临时目录:优先 /tmp(容器内),回退 /data,最后系统临时目录。"""
701
+ for candidate in ("/tmp/backup_v6_zip", "/data/backup_v6_zip"):
702
+ try:
703
+ p = Path(candidate)
704
+ p.mkdir(parents=True, exist_ok=True)
705
+ (p / ".wtest").touch(exist_ok=True)
706
+ (p / ".wtest").unlink(missing_ok=True)
707
+ return p
708
+ except Exception:
709
+ continue
710
+ import tempfile
711
+ p = Path(tempfile.gettempdir()) / "backup_v6_zip"
712
+ p.mkdir(parents=True, exist_ok=True)
713
+ return p
714
+
715
+
716
+ def build_volume_zip(set_id: str, volume: int) -> tuple[Path, int]:
717
+ """构建第 volume 卷 ZIP(含完整 manifest.json + 本卷 blob_<idx>.dat)。
718
+
719
+ 扁平结构:小天才 zip.decompress 对子目录 ZIP 不可靠(v5 实测),
720
+ blob 文件统一命名 blob_<idx>.dat,idx 与 manifest.blobs 数组下标对应。
721
+ 每卷都写一份 manifest.json,客户端任一卷即可拿到完整清单。
722
+ """
723
+ info = archive_info(set_id)
724
+ if volume < 0 or volume >= info["volumes"]:
725
+ raise ValueError("volume_out_of_range")
726
+ manifest = get_set_manifest(set_id)
727
+ blobs = manifest.get("blobs", []) or []
728
+
729
+ # 带 idx 的 manifest 副本
730
+ blobs_out = []
731
+ for i, entry in enumerate(blobs):
732
+ e = dict(entry) if isinstance(entry, dict) else {}
733
+ e["idx"] = i
734
+ blobs_out.append(e)
735
+ manifest_out = dict(manifest)
736
+ manifest_out["blobs"] = blobs_out
737
+ manifest_out["v6Volumes"] = info["volumes"]
738
+ manifest_out["v6CurrentVolume"] = volume
739
+
740
+ tmp_root = _zip_tmp_root()
741
+ # 唯一文件名:并发请求(双设备恢复同一备份点 / 重试重叠)互不覆盖,
742
+ # 各自的流式传输结束时删各自的临时文件
743
+ import uuid as _uuid
744
+ tmp_path = tmp_root / f"{set_id}_v{volume}_{_uuid.uuid4().hex[:8]}.zip"
745
+
746
+ entry_spec = info["entries"][volume]
747
+ with zipfile.ZipFile(tmp_path, "w", zipfile.ZIP_DEFLATED) as zf:
748
+ zf.writestr("manifest.json", json.dumps(manifest_out, ensure_ascii=False))
749
+ for i in range(entry_spec["startIdx"], entry_spec["endIdx"]):
750
+ b = blobs[i]
751
+ if not isinstance(b, dict):
752
+ continue
753
+ key = str(b.get("key") or "")
754
+ if not key:
755
+ continue
756
+ blob_path = get_blob_local_path(key)
757
+ if not blob_path or not blob_path.exists():
758
+ zf.writestr(f"blob_{i}.dat.MISSING", f"blob {key} missing")
759
+ continue
760
+ zf.write(blob_path, f"blob_{i}.dat")
761
+
762
+ return tmp_path, tmp_path.stat().st_size
763
+
764
+
765
+ def zip_iter_chunks(zip_path: Path, chunk_size: int = 64 * 1024):
766
+ """流式读取 ZIP 文件(传完后删除临时文件)。"""
767
+ async def _gen():
768
+ try:
769
+ with open(zip_path, "rb") as f:
770
+ while True:
771
+ chunk = f.read(chunk_size)
772
+ if not chunk:
773
+ break
774
+ yield chunk
775
+ except Exception as e:
776
+ logger.error("[v6_zip] iter chunks failed: %s", e)
777
+ raise
778
+ finally:
779
+ try:
780
+ zip_path.unlink(missing_ok=True)
781
+ except Exception:
782
+ pass
783
+ return _gen()
784
+
785
+
786
+ # ===== GC(垃圾回收)=====
787
+
788
+ def gc_orphan_blobs() -> dict:
789
+ """清理 refs=0 且上传超过宽限期的 blob(本地 + Hub + DB)。"""
790
+ cutoff = _now() - GC_BLOB_GRACE_SEC
791
+ with get_conn() as conn:
792
+ rows = conn.execute(
793
+ "SELECT key, size FROM backup_v6_blobs WHERE refs <= 0 AND uploaded_at < ?",
794
+ (cutoff,),
795
+ ).fetchall()
796
+ freed = 0
797
+ freed_bytes = 0
798
+ hub_deletes: list[str] = []
799
+ for r in rows:
800
+ key = str(r["key"])
801
+ freed += 1
802
+ freed_bytes += int(r["size"])
803
+ try:
804
+ with get_conn() as conn:
805
+ conn.execute("DELETE FROM backup_v6_blobs WHERE key = ?", (key,))
806
+ p = _blob_local_path(key)
807
+ if p.exists():
808
+ p.unlink(missing_ok=True)
809
+ try:
810
+ if p.parent.exists() and not any(p.parent.iterdir()):
811
+ p.parent.rmdir()
812
+ except OSError:
813
+ pass
814
+ if is_hub_enabled():
815
+ hub_deletes.append(key)
816
+ except Exception as e:
817
+ logger.warning("[v6_gc] delete blob failed key=%s: %s", key[:8], e)
818
+ if hub_deletes:
819
+ def _del_hub():
820
+ for key in hub_deletes:
821
+ try:
822
+ delete_from_hub(NS_BLOBS, f"{key[:2]}/{key}")
823
+ except Exception:
824
+ pass
825
+ try:
826
+ loop = asyncio.get_running_loop()
827
+ loop.run_in_executor(None, _del_hub)
828
+ except RuntimeError:
829
+ _del_hub()
830
+ return {"freedBlobs": freed, "freedBytes": freed_bytes}
831
+
832
+
833
+ _gc_task: "Optional[asyncio.Task]" = None
834
+
835
+
836
+ def start_gc_task() -> None:
837
+ global _gc_task
838
+ if _gc_task is not None:
839
+ return
840
+
841
+ async def _loop():
842
+ while True:
843
+ await asyncio.sleep(GC_INTERVAL_SEC)
844
+ try:
845
+ result = await asyncio.to_thread(gc_orphan_blobs)
846
+ if result["freedBlobs"] > 0:
847
+ logger.info("[v6_gc] freed %d blobs (%d bytes)",
848
+ result["freedBlobs"], result["freedBytes"])
849
+ except Exception as e:
850
+ logger.warning("[v6_gc] failed: %s", e)
851
+
852
+ try:
853
+ _gc_task = asyncio.create_task(_loop())
854
+ except RuntimeError:
855
+ pass
856
+
857
+
858
+ async def stop_gc_task() -> None:
859
+ global _gc_task
860
+ if _gc_task is not None:
861
+ _gc_task.cancel()
862
+ try:
863
+ await _gc_task
864
+ except asyncio.CancelledError:
865
+ pass
866
+ _gc_task = None
867
+
868
+
869
+ # ===== Hub 同步(HF Space 重建后恢复备份点元数据)=====
870
+
871
+ def sync_sets_from_hub() -> int:
872
+ """从 Hub 的 backup_v6_sets 命名空间拉回所有 manifest,重建 DB 记录。
873
+
874
+ HF Space 重建后本地 DB 丢失,但 Hub 上 manifest 仍在。
875
+ blob 内容本体按需从 Hub 拉回(get_blob_local_path 内置回源)。
876
+ 返回同步的 set 数。
877
+ """
878
+ if not is_hub_enabled():
879
+ return 0
880
+ synced = 0
881
+ try:
882
+ keys = list_hub_keys(NS_SETS)
883
+ except Exception as e:
884
+ logger.warning("[v6_sync] list hub keys failed: %s", e)
885
+ return 0
886
+ for sub_key in keys:
887
+ # sub_key 形如 <user_id>/<set_id>.json
888
+ parts = str(sub_key).split("/")
889
+ if len(parts) != 2 or not parts[0] or not parts[1].endswith(".json"):
890
+ continue
891
+ user_id, filename = parts[0], parts[1]
892
+ set_id = filename[:-5]
893
+ if not set_id:
894
+ continue
895
+ try:
896
+ with get_conn() as conn:
897
+ exists = conn.execute(
898
+ "SELECT 1 FROM backup_v6_sets WHERE set_id = ?", (set_id,)
899
+ ).fetchone()
900
+ if exists:
901
+ continue
902
+ data = download_bytes(NS_SETS, f"{user_id}/{set_id}.json")
903
+ if not data:
904
+ continue
905
+ manifest = json.loads(data.decode("utf-8"))
906
+ blobs = manifest.get("blobs", []) or []
907
+ keys_in_set = [str(b.get("key") or "") for b in blobs if isinstance(b, dict)]
908
+ total_size = sum(int(b.get("size") or 0) for b in blobs)
909
+ checksum = str(manifest.get("checksum") or "")
910
+ created_at = int(manifest.get("exportedAt") or _now())
911
+ # blob DB 记录补插(本地/Hub 有内容即可,refs 由各 set 累加)
912
+ with get_conn() as conn:
913
+ for i, k in enumerate(dict.fromkeys(keys_in_set)):
914
+ if not k:
915
+ continue
916
+ size = 0
917
+ mime = ""
918
+ for b in blobs:
919
+ if isinstance(b, dict) and str(b.get("key") or "") == k:
920
+ size = int(b.get("size") or 0)
921
+ mime = _guess_mime(str(b.get("path") or ""))
922
+ break
923
+ conn.execute(
924
+ "INSERT OR IGNORE INTO backup_v6_blobs(key, size, mime, refs, uploaded_at) "
925
+ "VALUES(?, ?, ?, 0, ?)",
926
+ (k, size, mime, created_at),
927
+ )
928
+ conn.execute(
929
+ "INSERT OR IGNORE INTO backup_v6_sets(set_id, user_id, alias, manifest, "
930
+ "manifest_sha256, total_size, blob_count, created_at, is_latest, prev_set_id, mode, device) "
931
+ "VALUES(?, ?, ?, ?, ?, ?, ?, ?, 0, ?, ?, ?)",
932
+ (
933
+ set_id, user_id, str(manifest.get("alias") or "")[:200],
934
+ json.dumps(manifest, ensure_ascii=False), checksum,
935
+ total_size, len(blobs), created_at,
936
+ str(manifest.get("prevSetId") or ""),
937
+ str(manifest.get("mode") or "")[:32],
938
+ str(manifest.get("device") or "")[:128],
939
+ ),
940
+ )
941
+ for k in dict.fromkeys(keys_in_set):
942
+ if k:
943
+ conn.execute(
944
+ "UPDATE backup_v6_blobs SET refs = refs + 1 WHERE key = ?", (k,)
945
+ )
946
+ # 重建 is_latest(最新 created_at 的 set)
947
+ with get_conn() as conn:
948
+ newest = conn.execute(
949
+ "SELECT set_id FROM backup_v6_sets WHERE user_id = ? ORDER BY created_at DESC LIMIT 1",
950
+ (user_id,),
951
+ ).fetchone()
952
+ conn.execute(
953
+ "UPDATE backup_v6_sets SET is_latest = 0 WHERE user_id = ?", (user_id,)
954
+ )
955
+ if newest:
956
+ conn.execute(
957
+ "UPDATE backup_v6_sets SET is_latest = 1 WHERE set_id = ?",
958
+ (str(newest["set_id"]),),
959
+ )
960
+ _write_manifest_to_local(user_id, set_id, manifest)
961
+ synced += 1
962
+ except Exception as e:
963
+ logger.warning("[v6_sync] sync set %s failed: %s", set_id[:8], e)
964
+ if synced:
965
+ logger.info("[v6_sync] restored %d sets from Hub", synced)
966
+ return synced
967
+
968
+
969
+ async def _sync_sets_background() -> None:
970
+ try:
971
+ await asyncio.to_thread(sync_sets_from_hub)
972
+ except Exception as e:
973
+ logger.warning("[v6_sync] background sync failed: %s", e)
974
+
975
+
976
+ # ===== v5 遗留数据清理 =====
977
+
978
+ def purge_legacy_v5(user_id: Optional[str] = None) -> dict:
979
+ """清理 v5 遗留数据(已被污染、设计缺陷无法正常使用)。
980
+
981
+ - user_id 指定:只清该用户的 v5 备份集 + 全部 refs=0 的 v5 blob + Hub 上该用户的 manifest
982
+ - user_id 为 None:清空全部 v5 数据(表 + 本地目录 + Hub 命名空间)
983
+
984
+ ★ 本函数会在工作线程(asyncio.to_thread)中执行,绝不调用任何含
985
+ asyncio.create_task / ensure_future 的函数(v5 的 delete_set 会在末尾
986
+ create_task 触发 GC,线程内无事件循环直接 RuntimeError),
987
+ 全部用同步 SQL + 文件操作实现。
988
+
989
+ 返回 {deletedSets, deletedBlobs, freedBytes}。
990
+ """
991
+ from ..hf_storage import _local_path
992
+
993
+ deleted_sets = 0
994
+ deleted_blobs = 0
995
+ freed_bytes = 0
996
+
997
+ def _v5_blob_path(key: str) -> Path:
998
+ return _local_path("backup_blobs_v5", f"{key[:2]}/{key}")
999
+
1000
+ if user_id:
1001
+ # 用户级:直接 SQL 批量删(不走 v5 delete_set,避免其内部 create_task)
1002
+ with get_conn() as conn:
1003
+ rows = conn.execute(
1004
+ "SELECT set_id, manifest FROM backup_sets WHERE user_id = ?", (user_id,)
1005
+ ).fetchall()
1006
+ for r in rows:
1007
+ set_id = str(r["set_id"])
1008
+ # 解析 manifest 拿 blob keys,直接减 refs(含 DB/本地/Hub 三处清理)
1009
+ keys: list[str] = []
1010
+ try:
1011
+ m = json.loads(r["manifest"])
1012
+ keys = [str(b.get("key") or "") for b in m.get("blobs", []) if isinstance(b, dict)]
1013
+ except Exception:
1014
+ keys = []
1015
+ try:
1016
+ with get_conn() as conn:
1017
+ for k in dict.fromkeys(keys):
1018
+ conn.execute(
1019
+ "UPDATE backup_blobs SET refs = MAX(refs - 1, 0) WHERE key = ?", (k,)
1020
+ )
1021
+ conn.execute("DELETE FROM backup_sets WHERE set_id = ?", (set_id,))
1022
+ deleted_sets += 1
1023
+ except Exception as e:
1024
+ logger.warning("[v6_purge] delete v5 set failed set=%s: %s", set_id[:8], e)
1025
+ # 删本地 + Hub 的 manifest
1026
+ try:
1027
+ mp = _local_path("backup_sets_v5", f"{user_id}/{set_id}.json")
1028
+ if mp.exists():
1029
+ mp.unlink(missing_ok=True)
1030
+ if is_hub_enabled():
1031
+ delete_from_hub("backup_sets_v5", f"{user_id}/{set_id}.json")
1032
+ except Exception:
1033
+ pass
1034
+ # 清理全部 refs=0 的 v5 blob(无论宽限期——v5 整体废弃,同步删)
1035
+ with get_conn() as conn:
1036
+ orphans = conn.execute(
1037
+ "SELECT key, size FROM backup_blobs WHERE refs <= 0"
1038
+ ).fetchall()
1039
+ for r in orphans:
1040
+ key = str(r["key"])
1041
+ try:
1042
+ with get_conn() as conn:
1043
+ conn.execute("DELETE FROM backup_blobs WHERE key = ?", (key,))
1044
+ conn.execute("DELETE FROM backup_blob_uploaders WHERE key = ?", (key,))
1045
+ p = _v5_blob_path(key)
1046
+ if p.exists():
1047
+ p.unlink(missing_ok=True)
1048
+ if is_hub_enabled():
1049
+ try:
1050
+ delete_from_hub("backup_blobs_v5", f"{key[:2]}/{key}")
1051
+ except Exception:
1052
+ pass
1053
+ deleted_blobs += 1
1054
+ freed_bytes += int(r["size"])
1055
+ except Exception as e:
1056
+ logger.warning("[v6_purge] delete v5 blob failed key=%s: %s", key[:8], e)
1057
+ # 删 Hub 上该用户的其余 v5 manifest(防御本地 DB 已缺记录的情况)
1058
+ if is_hub_enabled():
1059
+ try:
1060
+ for sub in list_hub_keys("backup_sets_v5"):
1061
+ if str(sub).split("/")[0] == user_id:
1062
+ try:
1063
+ delete_from_hub("backup_sets_v5", sub)
1064
+ except Exception:
1065
+ pass
1066
+ except Exception:
1067
+ pass
1068
+ else:
1069
+ # 全局清理:清表 + 清本地目录 + 清 Hub 命名空间
1070
+ with get_conn() as conn:
1071
+ rows = conn.execute("SELECT COUNT(*) AS c FROM backup_sets").fetchone()
1072
+ deleted_sets = int(rows["c"]) if rows else 0
1073
+ conn.execute("DELETE FROM backup_sets")
1074
+ conn.execute("DELETE FROM backup_blobs")
1075
+ conn.execute("DELETE FROM backup_blob_uploaders")
1076
+ # 本地文件
1077
+ import shutil
1078
+ for ns in ("backup_blobs_v5", "backup_sets_v5"):
1079
+ try:
1080
+ root = _local_path(ns, "_").parent
1081
+ if root.exists():
1082
+ shutil.rmtree(root, ignore_errors=True)
1083
+ except Exception:
1084
+ pass
1085
+ # Hub 命名空间
1086
+ if is_hub_enabled():
1087
+ for ns in ("backup_blobs_v5", "backup_sets_v5"):
1088
+ try:
1089
+ for sub in list_hub_keys(ns):
1090
+ try:
1091
+ delete_from_hub(ns, sub)
1092
+ except Exception:
1093
+ pass
1094
+ except Exception:
1095
+ pass
1096
+ deleted_blobs = -1 # 全局模式下不精确统计
1097
+
1098
+ return {"deletedSets": deleted_sets, "deletedBlobs": deleted_blobs, "freedBytes": freed_bytes}