diff --git a/main.py b/main.py index 0f535d7..5f7fb42 100644 --- a/main.py +++ b/main.py @@ -1,3 +1,6 @@ +import asyncio +from contextlib import asynccontextmanager, suppress + from fastapi import FastAPI from src.router import v1_router from src.utils.draft_downloader import download_draft @@ -5,8 +8,21 @@ from src.utils.logger import logger from src.middlewares import PrepareMiddleware, ResponseMiddleware, TraceContextMiddleware +@asynccontextmanager +async def lifespan(app: FastAPI): + from src.utils.draft_cleanup import draft_cleanup_background_loop + + cleanup_task = asyncio.create_task(draft_cleanup_background_loop()) + try: + yield + finally: + cleanup_task.cancel() + with suppress(asyncio.CancelledError): + await cleanup_task + + # 1. 创建 FastAPI 应用 -app: FastAPI = FastAPI(title="CapCut Mate API", version="1.0") +app: FastAPI = FastAPI(title="CapCut Mate API", version="1.0", lifespan=lifespan) # 2. 注册路由 app.include_router(router=v1_router, prefix="/openapi/capcut-mate", tags=["capcut-mate"]) diff --git a/src/utils/draft_cleanup.py b/src/utils/draft_cleanup.py new file mode 100644 index 0000000..3db3648 --- /dev/null +++ b/src/utils/draft_cleanup.py @@ -0,0 +1,150 @@ +""" +草稿目录定期清理:超出数量上限时按创建时间从旧到新删除目录,并跳过受保护、 +已加锁、仍在内存缓存中的草稿,避免影响正常读写。 +""" +from __future__ import annotations + +import asyncio +import datetime +import os +import re +import shutil +from typing import Iterable, Optional + +import config +from src.utils.draft_cache import DRAFT_CACHE +from src.utils.draft_lock_manager import get_draft_lock_manager +from src.utils.logger import logger + +# 草稿清理策略(内置常量,不放在 config) +DRAFT_CLEANUP_INTERVAL_SECONDS = 3600 +DRAFT_CLEANUP_MAX_DRAFT_COUNT = 1000 +DRAFT_CLEANUP_PROTECTED_DRAFT_IDS = frozenset( + { + "20251204214904ccb1af38", + "2025120421372636a27729", + } +) + +_DRAFT_ID_RE = re.compile(r"^\d{14}[a-f0-9]{8}$") + + +def is_draft_directory_name(name: str) -> bool: + return bool(_DRAFT_ID_RE.match(name)) + + +def _draft_sort_key(draft_id: str) -> tuple[datetime.datetime, str]: + ts = draft_id[:14] + dt = datetime.datetime.strptime(ts, "%Y%m%d%H%M%S") + return dt, draft_id + + +def list_sorted_draft_ids(draft_dir: str) -> list[str]: + if not os.path.isdir(draft_dir): + return [] + names: list[str] = [] + with os.scandir(draft_dir) as it: + for ent in it: + if ent.is_dir() and is_draft_directory_name(ent.name): + names.append(ent.name) + names.sort(key=_draft_sort_key) + return names + + +def select_drafts_for_deletion( + sorted_oldest_first: list[str], max_keep: int, skip_ids: set[str] +) -> list[str]: + n = len(sorted_oldest_first) + need = max(0, n - max_keep) + if need == 0: + return [] + out: list[str] = [] + for did in sorted_oldest_first: + if len(out) >= need: + break + if did in skip_ids: + continue + out.append(did) + return out + + +def delete_draft_folders(draft_dir: str, draft_ids: Iterable[str]) -> list[str]: + deleted: list[str] = [] + for did in draft_ids: + path = os.path.join(draft_dir, did) + logger.info( + "DRAFT_CLEANUP_DELETE draft_id=%s path=%s action=shutil.rmtree", + did, + path, + ) + try: + shutil.rmtree(path, ignore_errors=False) + deleted.append(did) + logger.info( + "DRAFT_CLEANUP_DELETED draft_id=%s path=%s status=success", + did, + path, + ) + except OSError as e: + logger.warning( + "DRAFT_CLEANUP_DELETE_FAILED draft_id=%s path=%s error=%s", + did, + path, + e, + ) + return deleted + + +def run_one_draft_cleanup( + draft_dir: Optional[str] = None, + max_keep: Optional[int] = None, + protected_ids: Optional[Iterable[str]] = None, + locked_ids: Optional[Iterable[str]] = None, + cached_ids: Optional[Iterable[str]] = None, +) -> list[str]: + base = draft_dir if draft_dir is not None else config.DRAFT_DIR + limit = max_keep if max_keep is not None else DRAFT_CLEANUP_MAX_DRAFT_COUNT + protected = frozenset( + protected_ids + if protected_ids is not None + else DRAFT_CLEANUP_PROTECTED_DRAFT_IDS + ) + locked = ( + set(locked_ids) + if locked_ids is not None + else set(get_draft_lock_manager().get_all_locked_drafts()) + ) + cached = set(cached_ids) if cached_ids is not None else set(DRAFT_CACHE.keys()) + skip = set(protected) | locked | cached + sorted_ids = list_sorted_draft_ids(base) + to_remove = select_drafts_for_deletion(sorted_ids, limit, skip) + if not to_remove: + reason = "under_limit" if len(sorted_ids) <= limit else "all_excess_protected_or_active" + logger.info( + "DRAFT_CLEANUP_SKIP total_drafts=%s max_keep=%s skip_count=%s reason=%s", + len(sorted_ids), + limit, + len(skip), + reason, + ) + return [] + logger.info( + "DRAFT_CLEANUP_PLAN total_drafts=%s max_keep=%s to_delete=%s protected=%s locked=%s cached=%s", + len(sorted_ids), + limit, + len(to_remove), + len(protected), + len(locked), + len(cached), + ) + return delete_draft_folders(base, to_remove) + + +async def draft_cleanup_background_loop() -> None: + interval = DRAFT_CLEANUP_INTERVAL_SECONDS + while True: + try: + run_one_draft_cleanup() + except Exception: + logger.exception("DRAFT_CLEANUP_ERROR phase=run_one_draft_cleanup") + await asyncio.sleep(interval) diff --git a/tests/test_draft_cleanup.py b/tests/test_draft_cleanup.py new file mode 100644 index 0000000..9981251 --- /dev/null +++ b/tests/test_draft_cleanup.py @@ -0,0 +1,145 @@ +"""Unit tests for draft directory cleanup (oldest first, protected / locked / cached skips).""" +from __future__ import annotations + +import os + +from src.utils import draft_cleanup as dc + + +def _mkdir(p: str, name: str) -> str: + path = os.path.join(p, name) + os.makedirs(path, exist_ok=True) + return path + + +def test_is_draft_directory_name_accepts_standard_id() -> None: + assert dc.is_draft_directory_name("20200101000000aaaaaaaa") is True + assert dc.is_draft_directory_name("20200101000000abcdef12") is True + + +def test_is_draft_directory_name_rejects_invalid() -> None: + assert dc.is_draft_directory_name("20200101000000aaaaaaa") is False + assert dc.is_draft_directory_name("20200101000000aaaaaaaag") is False + assert dc.is_draft_directory_name("draft") is False + assert dc.is_draft_directory_name("") is False + + +def test_list_sorted_draft_ids_ignores_files_and_bad_names(tmp_path) -> None: + base = str(tmp_path) + _mkdir(base, "20200201000000bbbbbbbb") + _mkdir(base, "20200101000000aaaaaaaa") + open(os.path.join(base, "x.mp4"), "wb").close() + _mkdir(base, "not_a_valid_draft_id") + got = dc.list_sorted_draft_ids(base) + assert got == ["20200101000000aaaaaaaa", "20200201000000bbbbbbbb"] + + +def test_select_drafts_for_deletion_oldest_first_until_quota(tmp_path) -> None: + ordered = [ + "20200101000000aaaaaaaa", + "20200201000000bbbbbbbb", + "20200301000000cccccccc", + ] + assert dc.select_drafts_for_deletion(ordered, max_keep=3, skip_ids=set()) == [] + assert dc.select_drafts_for_deletion(ordered, max_keep=2, skip_ids=set()) == [ + "20200101000000aaaaaaaa" + ] + assert dc.select_drafts_for_deletion(ordered, max_keep=1, skip_ids=set()) == [ + "20200101000000aaaaaaaa", + "20200201000000bbbbbbbb", + ] + assert dc.select_drafts_for_deletion(ordered, max_keep=0, skip_ids=set()) == ordered + + +def test_select_drafts_for_deletion_skips_protected_uses_next_oldest() -> None: + ordered = [ + "20200101000000aaaaaaaa", + "20200201000000bbbbbbbb", + "20200301000000cccccccc", + ] + skip = {"20200101000000aaaaaaaa"} + assert dc.select_drafts_for_deletion(ordered, max_keep=1, skip_ids=skip) == [ + "20200201000000bbbbbbbb", + "20200301000000cccccccc", + ] + + +def test_run_one_draft_cleanup_deletes_oldest_unskipped_only(tmp_path) -> None: + base = str(tmp_path) + old = "20200101000000aaaaaaaa" + mid = "20200201000000bbbbbbbb" + new = "20200301000000cccccccc" + for d in (old, mid, new): + _mkdir(base, d) + deleted = dc.run_one_draft_cleanup( + draft_dir=base, + max_keep=1, + protected_ids=[], + locked_ids=[], + cached_ids=[], + ) + assert set(deleted) == {old, mid} + assert os.path.isdir(os.path.join(base, new)) + assert not os.path.exists(os.path.join(base, old)) + assert not os.path.exists(os.path.join(base, mid)) + + +def test_run_one_draft_cleanup_never_removes_protected_ids(tmp_path) -> None: + base = str(tmp_path) + for pid in dc.DRAFT_CLEANUP_PROTECTED_DRAFT_IDS: + _mkdir(base, pid) + extra = "20990101000000eeeeeeee" + _mkdir(base, extra) + deleted = dc.run_one_draft_cleanup( + draft_dir=base, + max_keep=0, + protected_ids=dc.DRAFT_CLEANUP_PROTECTED_DRAFT_IDS, + locked_ids=[], + cached_ids=[], + ) + assert extra in deleted + for pid in dc.DRAFT_CLEANUP_PROTECTED_DRAFT_IDS: + assert pid not in deleted + assert os.path.isdir(os.path.join(base, pid)) + + +def test_run_one_draft_cleanup_skips_locked_ids(tmp_path) -> None: + base = str(tmp_path) + a, b, c = ( + "20200101000000aaaaaaaa", + "20200201000000bbbbbbbb", + "20200301000000cccccccc", + ) + for d in (a, b, c): + _mkdir(base, d) + deleted = dc.run_one_draft_cleanup( + draft_dir=base, + max_keep=1, + protected_ids=[], + locked_ids=[a], + cached_ids=[], + ) + assert a not in deleted + assert set(deleted) == {b, c} + assert os.path.isdir(os.path.join(base, a)) + + +def test_run_one_draft_cleanup_skips_cached_ids(tmp_path) -> None: + base = str(tmp_path) + a, b = "20200101000000aaaaaaaa", "20200201000000bbbbbbbb" + for d in (a, b): + _mkdir(base, d) + deleted = dc.run_one_draft_cleanup( + draft_dir=base, + max_keep=1, + protected_ids=[], + locked_ids=[], + cached_ids=[a], + ) + assert a not in deleted + assert deleted == [b] + + +def test_list_sorted_draft_ids_missing_dir_returns_empty(tmp_path) -> None: + missing = os.path.join(str(tmp_path), "nope") + assert dc.list_sorted_draft_ids(missing) == []