library_trash.py 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489
  1. """Library trash sweeper + purge service (#1008).
  2. Two-stage file deletion for the library:
  3. 1. Users / admins soft-delete files — the row stays in ``library_files`` with
  4. ``deleted_at`` stamped; the bytes stay on disk. This is handled inline in
  5. ``backend.app.api.routes.library`` and exposed to admins as a bulk "purge
  6. old files" operation via :meth:`LibraryTrashService.purge_older_than`.
  7. 2. A background sweeper in this service hard-deletes rows (and their bytes)
  8. whose ``deleted_at`` is older than the configured retention window.
  9. External files (``is_external=True``) are never placed in the trash — their
  10. bytes live outside Bambuddy's control, so there's nothing to restore.
  11. """
  12. from __future__ import annotations
  13. import asyncio
  14. import logging
  15. from datetime import datetime, timedelta, timezone
  16. from pathlib import Path
  17. from sqlalchemy import and_, delete, func, or_, select
  18. from sqlalchemy.ext.asyncio import AsyncSession
  19. from backend.app.core.config import settings as app_settings
  20. from backend.app.core.database import async_session
  21. from backend.app.models.library import LibraryFile
  22. from backend.app.models.print_queue import PrintQueueItem, PrintQueueVariant
  23. from backend.app.models.settings import Settings
  24. from backend.app.utils.library_paths import remove_library_photos_dir
  25. from backend.app.utils.local_time import utcnow_naive
  26. logger = logging.getLogger(__name__)
  27. # Settings key used to persist the trash retention window (days). The sweeper
  28. # reads this on every tick so the UI can change it without a restart.
  29. TRASH_RETENTION_KEY = "library_trash_retention_days"
  30. DEFAULT_RETENTION_DAYS = 30
  31. # Clamp retention to a sensible range. 1 day is a reasonable floor (anything
  32. # shorter just makes trash into hard-delete); 365 gives admins plenty of rope
  33. # without letting accidental typos (99999) grow the table unboundedly.
  34. MIN_RETENTION_DAYS = 1
  35. MAX_RETENTION_DAYS = 365
  36. # Auto-purge settings (#1008 follow-up). When enabled, the sweeper loop also
  37. # runs the admin bulk purge once per 24h using the saved age threshold.
  38. # Default-off so existing installs don't surprise users — opt-in via Settings.
  39. AUTO_PURGE_ENABLED_KEY = "library_auto_purge_enabled"
  40. AUTO_PURGE_DAYS_KEY = "library_auto_purge_days"
  41. AUTO_PURGE_INCLUDE_NEVER_PRINTED_KEY = "library_auto_purge_include_never_printed"
  42. AUTO_PURGE_LAST_RUN_KEY = "library_auto_purge_last_run"
  43. DEFAULT_AUTO_PURGE_DAYS = 90
  44. MIN_AUTO_PURGE_DAYS = 7 # anything shorter is begging for accidents
  45. MAX_AUTO_PURGE_DAYS = 3650
  46. def _to_absolute_path(relative_path: str | None) -> Path | None:
  47. """Mirror of the routes helper so this service has no route-module import.
  48. Accepts the legacy absolute paths that predate the relative-path migration
  49. verbatim; new rows always store paths relative to ``base_dir``.
  50. """
  51. if not relative_path:
  52. return None
  53. path = Path(relative_path)
  54. if path.is_absolute():
  55. return path
  56. return (
  57. Path(app_settings.base_dir) / path
  58. ) # SEC-PATH-OK: relative_path is LibraryFile.file_path / LibraryFile.thumbnail_path — DB-stored, internally generated by the upload pipeline
  59. def _age_cutoff(now: datetime, older_than_days: int) -> datetime:
  60. return now - timedelta(days=older_than_days)
  61. def _purge_filter(cutoff: datetime, include_never_printed: bool):
  62. """SQLAlchemy clause selecting files eligible for admin purge.
  63. A file is "old" if either (a) ``last_printed_at`` is set and predates the
  64. cutoff, or (b) ``last_printed_at`` is NULL *and* the file was uploaded
  65. before the cutoff — but only when ``include_never_printed`` is True.
  66. """
  67. last_printed_old = and_(
  68. LibraryFile.last_printed_at.isnot(None),
  69. LibraryFile.last_printed_at < cutoff,
  70. )
  71. if include_never_printed:
  72. never_printed_old = and_(
  73. LibraryFile.last_printed_at.is_(None),
  74. LibraryFile.created_at < cutoff,
  75. )
  76. age_clause = or_(last_printed_old, never_printed_old)
  77. else:
  78. age_clause = last_printed_old
  79. return and_(
  80. LibraryFile.deleted_at.is_(None),
  81. LibraryFile.is_external.is_(False),
  82. age_clause,
  83. )
  84. class LibraryTrashService:
  85. """Manages the trash retention sweeper and admin-triggered bulk purges."""
  86. def __init__(self):
  87. self._scheduler_task: asyncio.Task | None = None
  88. # Tick every 15 minutes — the window is a day, so this is plenty
  89. # responsive without burning CPU.
  90. self._check_interval = 900
  91. async def start_scheduler(self):
  92. """Start the background sweeper task (idempotent)."""
  93. if self._scheduler_task is not None:
  94. return
  95. logger.info("Starting library trash sweeper")
  96. self._scheduler_task = asyncio.create_task(self._scheduler_loop())
  97. def stop_scheduler(self):
  98. if self._scheduler_task:
  99. self._scheduler_task.cancel()
  100. self._scheduler_task = None
  101. logger.info("Stopped library trash sweeper")
  102. async def _scheduler_loop(self):
  103. while True:
  104. try:
  105. await asyncio.sleep(self._check_interval)
  106. async with async_session() as db:
  107. await self._sweep(db)
  108. await self._maybe_run_auto_purge(db)
  109. except asyncio.CancelledError:
  110. break
  111. except Exception as e: # pragma: no cover - defensive
  112. logger.error("Error in library trash sweeper: %s", e)
  113. await asyncio.sleep(60)
  114. # ---- Settings -----------------------------------------------------
  115. async def get_retention_days(self, db: AsyncSession | None = None) -> int:
  116. if db is None:
  117. async with async_session() as session:
  118. return await self._read_retention(session)
  119. return await self._read_retention(db)
  120. @staticmethod
  121. async def _read_retention(db: AsyncSession) -> int:
  122. result = await db.execute(select(Settings.value).where(Settings.key == TRASH_RETENTION_KEY))
  123. raw = result.scalar_one_or_none()
  124. if raw is None:
  125. return DEFAULT_RETENTION_DAYS
  126. try:
  127. days = int(raw)
  128. except (TypeError, ValueError):
  129. return DEFAULT_RETENTION_DAYS
  130. return max(MIN_RETENTION_DAYS, min(MAX_RETENTION_DAYS, days))
  131. async def set_retention_days(self, db: AsyncSession, days: int) -> int:
  132. """Persist the retention window. Clamped to [MIN, MAX]."""
  133. clamped = max(MIN_RETENTION_DAYS, min(MAX_RETENTION_DAYS, int(days)))
  134. result = await db.execute(select(Settings).where(Settings.key == TRASH_RETENTION_KEY))
  135. row = result.scalar_one_or_none()
  136. if row is None:
  137. db.add(Settings(key=TRASH_RETENTION_KEY, value=str(clamped)))
  138. else:
  139. row.value = str(clamped)
  140. await db.commit()
  141. return clamped
  142. @staticmethod
  143. async def _read_setting(db: AsyncSession, key: str) -> str | None:
  144. result = await db.execute(select(Settings.value).where(Settings.key == key))
  145. return result.scalar_one_or_none()
  146. @staticmethod
  147. async def _write_setting(db: AsyncSession, key: str, value: str) -> None:
  148. result = await db.execute(select(Settings).where(Settings.key == key))
  149. row = result.scalar_one_or_none()
  150. if row is None:
  151. db.add(Settings(key=key, value=value))
  152. else:
  153. row.value = value
  154. async def get_auto_purge_settings(self, db: AsyncSession) -> dict:
  155. """Return the current auto-purge config.
  156. Returns a dict with ``enabled`` (bool), ``days`` (int, clamped) and
  157. ``include_never_printed`` (bool). Missing keys default to disabled /
  158. 90 days / include-never-printed-on, matching the manual purge UX.
  159. """
  160. enabled_raw = await self._read_setting(db, AUTO_PURGE_ENABLED_KEY)
  161. days_raw = await self._read_setting(db, AUTO_PURGE_DAYS_KEY)
  162. incl_raw = await self._read_setting(db, AUTO_PURGE_INCLUDE_NEVER_PRINTED_KEY)
  163. enabled = (enabled_raw or "false").lower() == "true"
  164. try:
  165. days = int(days_raw) if days_raw is not None else DEFAULT_AUTO_PURGE_DAYS
  166. except (TypeError, ValueError):
  167. days = DEFAULT_AUTO_PURGE_DAYS
  168. days = max(MIN_AUTO_PURGE_DAYS, min(MAX_AUTO_PURGE_DAYS, days))
  169. include_never_printed = (incl_raw or "true").lower() == "true"
  170. return {
  171. "enabled": enabled,
  172. "days": days,
  173. "include_never_printed": include_never_printed,
  174. }
  175. async def set_auto_purge_settings(
  176. self,
  177. db: AsyncSession,
  178. *,
  179. enabled: bool,
  180. days: int,
  181. include_never_printed: bool,
  182. ) -> dict:
  183. """Persist auto-purge config; returns the saved (clamped) values."""
  184. clamped_days = max(MIN_AUTO_PURGE_DAYS, min(MAX_AUTO_PURGE_DAYS, int(days)))
  185. await self._write_setting(db, AUTO_PURGE_ENABLED_KEY, "true" if enabled else "false")
  186. await self._write_setting(db, AUTO_PURGE_DAYS_KEY, str(clamped_days))
  187. await self._write_setting(
  188. db,
  189. AUTO_PURGE_INCLUDE_NEVER_PRINTED_KEY,
  190. "true" if include_never_printed else "false",
  191. )
  192. await db.commit()
  193. return {
  194. "enabled": enabled,
  195. "days": clamped_days,
  196. "include_never_printed": include_never_printed,
  197. }
  198. async def _get_last_auto_purge_run(self, db: AsyncSession) -> datetime | None:
  199. raw = await self._read_setting(db, AUTO_PURGE_LAST_RUN_KEY)
  200. if not raw:
  201. return None
  202. try:
  203. # Stored as ISO 8601 UTC; tolerate both with and without 'Z' suffix.
  204. return datetime.fromisoformat(raw.replace("Z", "+00:00"))
  205. except ValueError:
  206. return None
  207. async def _stamp_last_auto_purge_run(self, db: AsyncSession, when: datetime) -> None:
  208. await self._write_setting(db, AUTO_PURGE_LAST_RUN_KEY, when.isoformat())
  209. await db.commit()
  210. async def _maybe_run_auto_purge(self, db: AsyncSession) -> int:
  211. """If auto-purge is enabled and >=24h has elapsed since the last run, run it.
  212. Returns the number of files moved to trash (0 if disabled or throttled).
  213. The 24h throttle means a 15-minute sweeper cadence still only triggers
  214. one actual purge per day, keeping the DB churn predictable.
  215. """
  216. cfg = await self.get_auto_purge_settings(db)
  217. if not cfg["enabled"]:
  218. return 0
  219. now = datetime.now(timezone.utc)
  220. last = await self._get_last_auto_purge_run(db)
  221. if last is not None and (now - last) < timedelta(hours=24):
  222. return 0
  223. moved = await self.purge_older_than(
  224. db,
  225. older_than_days=cfg["days"],
  226. include_never_printed=cfg["include_never_printed"],
  227. )
  228. await self._stamp_last_auto_purge_run(db, now)
  229. if moved:
  230. logger.info("Library auto-purge: moved %d file(s) to trash (threshold=%d days)", moved, cfg["days"])
  231. return moved
  232. # ---- Preview / purge ---------------------------------------------
  233. async def preview_purge(
  234. self,
  235. db: AsyncSession,
  236. older_than_days: int,
  237. include_never_printed: bool = True,
  238. sample_limit: int = 5,
  239. ) -> dict:
  240. """Count + size of files eligible for purge. Reads only; never mutates."""
  241. if older_than_days < 1:
  242. return {"count": 0, "total_bytes": 0, "sample_filenames": []}
  243. now = datetime.now(timezone.utc)
  244. cutoff = _age_cutoff(now, older_than_days)
  245. clause = _purge_filter(cutoff, include_never_printed)
  246. count_result = await db.execute(select(func.count(LibraryFile.id)).where(clause))
  247. count = int(count_result.scalar() or 0)
  248. size_result = await db.execute(select(func.coalesce(func.sum(LibraryFile.file_size), 0)).where(clause))
  249. total_bytes = int(size_result.scalar() or 0)
  250. sample_result = await db.execute(
  251. select(LibraryFile.filename).where(clause).order_by(LibraryFile.created_at).limit(sample_limit)
  252. )
  253. samples = [row[0] for row in sample_result.all()]
  254. return {
  255. "count": count,
  256. "total_bytes": total_bytes,
  257. "sample_filenames": samples,
  258. "older_than_days": older_than_days,
  259. "include_never_printed": include_never_printed,
  260. }
  261. async def purge_older_than(
  262. self,
  263. db: AsyncSession,
  264. older_than_days: int,
  265. include_never_printed: bool = True,
  266. ) -> int:
  267. """Move matching files to trash (stamps ``deleted_at``). Returns count."""
  268. if older_than_days < 1:
  269. return 0
  270. now = datetime.now(timezone.utc)
  271. cutoff = _age_cutoff(now, older_than_days)
  272. clause = _purge_filter(cutoff, include_never_printed)
  273. # We need the IDs so callers can audit or display them if they want.
  274. # Doing a single UPDATE ... WHERE is safe even under concurrent
  275. # uploads — the clause already excludes rows with deleted_at set.
  276. id_result = await db.execute(select(LibraryFile.id).where(clause))
  277. ids = [row[0] for row in id_result.all()]
  278. if not ids:
  279. return 0
  280. await db.execute(LibraryFile.__table__.update().where(LibraryFile.id.in_(ids)).values(deleted_at=now))
  281. await db.commit()
  282. logger.info("Library purge: moved %d file(s) to trash (older_than_days=%d)", len(ids), older_than_days)
  283. return len(ids)
  284. # ---- Sweeper ------------------------------------------------------
  285. async def _sweep(self, db: AsyncSession) -> int:
  286. """Hard-delete trashed rows whose retention window has elapsed."""
  287. retention = await self._read_retention(db)
  288. now = datetime.now(timezone.utc)
  289. cutoff = now - timedelta(days=retention)
  290. result = await db.execute(
  291. select(LibraryFile).where(
  292. LibraryFile.deleted_at.isnot(None),
  293. LibraryFile.deleted_at < cutoff,
  294. )
  295. )
  296. rows = result.scalars().all()
  297. if not rows:
  298. return 0
  299. deleted = 0
  300. for row in rows:
  301. self._unlink_on_disk(row)
  302. deleted += 1
  303. await delete_dependent_variants(db, [r.id for r in rows])
  304. await release_queue_references(db, [r.id for r in rows])
  305. # Single DELETE is faster than N await db.delete() round-trips; we
  306. # still need the Python loop above to unlink bytes on disk.
  307. await db.execute(delete(LibraryFile).where(LibraryFile.id.in_([r.id for r in rows])))
  308. await db.commit()
  309. logger.info("Library trash sweeper: hard-deleted %d row(s) past %d-day retention", deleted, retention)
  310. return deleted
  311. @staticmethod
  312. def _unlink_on_disk(row: LibraryFile) -> None:
  313. """Best-effort cleanup of the file, thumbnail and photos (#3077) on disk."""
  314. for rel in (row.file_path, row.thumbnail_path):
  315. abs_path = _to_absolute_path(rel)
  316. if abs_path is None:
  317. continue
  318. try:
  319. if abs_path.exists():
  320. abs_path.unlink()
  321. except OSError as e:
  322. logger.warning("Trash sweep: failed to unlink %s: %s", abs_path, e)
  323. remove_library_photos_dir(row.id)
  324. # ---- User-facing trash ops ----------------------------------------
  325. async def restore(self, db: AsyncSession, file: LibraryFile) -> LibraryFile:
  326. """Clear ``deleted_at`` so the file reappears in listings."""
  327. file.deleted_at = None
  328. await db.commit()
  329. await db.refresh(file)
  330. return file
  331. async def hard_delete_now(self, db: AsyncSession, file: LibraryFile) -> None:
  332. """Bypass retention and delete this trashed file + its bytes immediately."""
  333. self._unlink_on_disk(file)
  334. await delete_dependent_variants(db, [file.id])
  335. await release_queue_references(db, [file.id])
  336. await db.delete(file)
  337. await db.commit()
  338. async def release_queue_references(db: AsyncSession, file_ids: list[int]) -> int:
  339. """Take queued work off files that are about to be hard-deleted (#2819).
  340. Call this before any statement that removes ``library_files`` rows — the
  341. plain deletes in the routes, the folder cascade, and the sweeper. It is the
  342. same repair the scheduler does when a dispatch consumes its own library row
  343. (``_repoint_siblings_at_archive``), minus the part that cannot apply here:
  344. nothing is being printed, so there is no archive to hand the work to.
  345. Two things happen, and both matter on a different database:
  346. * Items still waiting on one of these files are cancelled, saying which file
  347. went. Without it a queued job sat there looking dispatchable and failed at
  348. the printer with "Library file not found", days later and with nothing
  349. naming the delete that caused it.
  350. * Every remaining row referencing the file has ``library_file_id`` cleared.
  351. That is what keeps it: ``print_queue.library_file_id`` is ``ON DELETE
  352. CASCADE``, which SQLite does not enforce and PostgreSQL does, so those rows
  353. were silently deleted there -- including finished ones, which is what a
  354. batch order counts its progress from.
  355. Rows already printing are left in place. One of those is a job on a machine
  356. right now; the file being deleted is the copy in the library, not the copy
  357. the printer is working from. Returns the number of items cancelled.
  358. """
  359. if not file_ids:
  360. return 0
  361. doomed: dict[int, list[int]] = {}
  362. rows = (
  363. await db.execute(
  364. select(PrintQueueItem.id, PrintQueueItem.library_file_id)
  365. .where(PrintQueueItem.library_file_id.in_(file_ids))
  366. .where(PrintQueueItem.archive_id.is_(None))
  367. # "skipped" is not terminal: clearing a printer's previous-success
  368. # gate puts those items back to pending, onto a file that by then
  369. # is gone.
  370. .where(PrintQueueItem.status.in_(("pending", "skipped")))
  371. )
  372. ).all()
  373. for item_id, lib_id in rows:
  374. doomed.setdefault(lib_id, []).append(item_id)
  375. if doomed:
  376. names = dict(
  377. (
  378. await db.execute(select(LibraryFile.id, LibraryFile.filename).where(LibraryFile.id.in_(list(doomed))))
  379. ).all()
  380. )
  381. # Naive UTC: `completed_at` is a naive column, and asyncpg rejects an
  382. # aware value outright where SQLite silently drops the offset.
  383. now = utcnow_naive()
  384. # One statement per file rather than per item: the case this exists for
  385. # is many copies of one file, and a folder delete can reach a lot of
  386. # them at once.
  387. for lib_id, item_ids in doomed.items():
  388. await db.execute(
  389. PrintQueueItem.__table__.update()
  390. .where(PrintQueueItem.id.in_(item_ids))
  391. .values(
  392. status="cancelled",
  393. completed_at=now,
  394. error_message=f"'{names.get(lib_id, 'The library file')}' was deleted from the library",
  395. )
  396. )
  397. logger.info("Library delete: cancelled %d queued item(s) whose file was removed", len(rows))
  398. await db.execute(
  399. PrintQueueItem.__table__.update()
  400. .where(PrintQueueItem.library_file_id.in_(file_ids))
  401. .values(library_file_id=None)
  402. )
  403. return len(rows)
  404. async def delete_dependent_variants(db: AsyncSession, file_ids: list[int]) -> None:
  405. """Drop cross-model queue candidates that pointed at these files (#671).
  406. SQLite ships with ``PRAGMA foreign_keys`` off — verified, not assumed — so
  407. the ON DELETE CASCADE on ``print_queue_variants.library_file_id`` never fires
  408. on the default deployment and the rows would outlive the file.
  409. The scheduler already refuses to dispatch a candidate whose file is missing
  410. or trashed, so nothing prints wrongly without this. It is here so the table
  411. does not fill with rows referencing files that no longer exist.
  412. """
  413. if not file_ids:
  414. return
  415. await db.execute(delete(PrintQueueVariant).where(PrintQueueVariant.library_file_id.in_(file_ids)))
  416. library_trash_service = LibraryTrashService()