"""Хранилище задач на SQLite. SQLite, а не память: задачи должны переживать перезапуск сервиса, иначе после обновления или падения клиент навсегда останется с id, по которому ничего нет. """ import json import sqlite3 import threading import time import uuid from pathlib import Path __all__ = ["JobStore", "JobStatus"] class JobStatus: QUEUED = "queued" RUNNING = "running" DONE = "done" FAILED = "failed" _SCHEMA = """ CREATE TABLE IF NOT EXISTS jobs ( id TEXT PRIMARY KEY, filename TEXT NOT NULL, duration_sec REAL NOT NULL, status TEXT NOT NULL, created_at REAL NOT NULL, started_at REAL, finished_at REAL, result TEXT, error TEXT, options TEXT ); CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs(status, created_at); """ class JobStore: def __init__(self, path: Path): self._path = Path(path) self._path.parent.mkdir(parents=True, exist_ok=True) # Один процесс, несколько потоков (веб + воркер), поэтому check_same_thread # выключен, а доступ сериализуется собственным замком. self._lock = threading.Lock() self._conn = sqlite3.connect(str(self._path), check_same_thread=False) self._conn.row_factory = sqlite3.Row self._conn.executescript(_SCHEMA) self._conn.commit() self._recover_stuck() self._backfill_duration() def _backfill_duration(self) -> None: """Проставляет длительность задачам, обработанным до её сохранения. Раньше поле оставалось нулевым: при загрузке файл ещё не декодирован, а после обработки никто не возвращал значение обратно. В сохранённом результате оно есть, поэтому историю можно починить. """ with self._lock: rows = self._conn.execute( "SELECT id, result FROM jobs " "WHERE duration_sec <= 0 AND result IS NOT NULL").fetchall() for row in rows: try: value = float(json.loads(row["result"]).get("duration_sec") or 0) except (ValueError, TypeError, json.JSONDecodeError): continue if value > 0: self._update(row["id"], duration_sec=value) def _recover_stuck(self) -> None: """После аварийного перезапуска задачи «в работе» никто не доделает.""" with self._lock: self._conn.execute( "UPDATE jobs SET status=?, error=? WHERE status=?", (JobStatus.FAILED, "сервис был перезапущен во время обработки", JobStatus.RUNNING), ) self._conn.commit() def create(self, filename: str, duration_sec: float, options: dict | None = None) -> str: job_id = uuid.uuid4().hex with self._lock: self._conn.execute( "INSERT INTO jobs (id, filename, duration_sec, status, created_at, options)" " VALUES (?,?,?,?,?,?)", (job_id, filename, duration_sec, JobStatus.QUEUED, time.time(), json.dumps(options or {}, ensure_ascii=False)), ) self._conn.commit() return job_id def get(self, job_id: str) -> dict | None: with self._lock: row = self._conn.execute("SELECT * FROM jobs WHERE id=?", (job_id,)).fetchone() if row is None: return None job = dict(row) job["result"] = json.loads(job["result"]) if job["result"] else None job["options"] = json.loads(job["options"]) if job["options"] else {} return job def claim_next(self) -> str | None: """Забирает самую старую задачу из очереди и сразу помечает её в работе. Выборка и пометка выполняются одним оператором под общим замком: иначе два воркера успевают увидеть одну и ту же задачу и берут её оба. """ with self._lock: row = self._conn.execute( "UPDATE jobs SET status=?, started_at=? WHERE id = (" " SELECT id FROM jobs WHERE status=? ORDER BY created_at LIMIT 1" ") RETURNING id", (JobStatus.RUNNING, time.time(), JobStatus.QUEUED), ).fetchone() self._conn.commit() return row["id"] if row else None def queue_position(self, job_id: str) -> int: """Сколько задач будет обработано до этой. Считается и та, что уже выполняется: для клиента она тоже стоит впереди, и ждать ему придётся вместе с ней. 0 означает «следующая на обработку». """ with self._lock: row = self._conn.execute( "SELECT COUNT(*) AS n FROM jobs WHERE status IN (?,?) AND created_at <" " (SELECT created_at FROM jobs WHERE id=?)", (JobStatus.QUEUED, JobStatus.RUNNING, job_id), ).fetchone() return row["n"] if row else 0 def mark_running(self, job_id: str) -> None: self._update(job_id, status=JobStatus.RUNNING, started_at=time.time()) def mark_done(self, job_id: str, result: dict) -> None: # Длительность становится известна только после декодирования, а в # списке задач нужна без разбора всего результата - записываем сюда. self._update(job_id, status=JobStatus.DONE, finished_at=time.time(), duration_sec=float(result.get("duration_sec") or 0.0), result=json.dumps(result, ensure_ascii=False)) def mark_failed(self, job_id: str, error: str) -> None: self._update(job_id, status=JobStatus.FAILED, finished_at=time.time(), error=error) def _update(self, job_id: str, **fields) -> None: assigns = ", ".join(f"{k}=?" for k in fields) with self._lock: self._conn.execute(f"UPDATE jobs SET {assigns} WHERE id=?", (*fields.values(), job_id)) self._conn.commit() def delete(self, job_id: str) -> bool: with self._lock: cur = self._conn.execute("DELETE FROM jobs WHERE id=?", (job_id,)) self._conn.commit() return cur.rowcount > 0 def cleanup(self, max_age_hours: float) -> list[str]: """Удаляет завершённые задачи старше срока. Возвращает их идентификаторы. Идентификаторы нужны вызывающему: рядом с записью в базе лежит файл разговора, и удалять их надо вместе, иначе папка растёт вечно. """ cutoff = time.time() - max_age_hours * 3600 with self._lock: rows = self._conn.execute( "SELECT id FROM jobs WHERE status IN (?,?) " "AND COALESCE(finished_at, created_at) <= ?", (JobStatus.DONE, JobStatus.FAILED, cutoff)).fetchall() ids = [row["id"] for row in rows] self._conn.execute( "DELETE FROM jobs WHERE status IN (?,?) AND COALESCE(finished_at, created_at) <= ?", (JobStatus.DONE, JobStatus.FAILED, cutoff), ) self._conn.commit() removed = len(ids) if removed: self.vacuum() return ids def vacuum(self) -> None: """Возвращает системе место, освобождённое удалением. SQLite помечает страницы свободными, но файл не уменьшает: после чистки многомесячной истории он остаётся прежнего размера. """ with self._lock: self._conn.execute("VACUUM") self._conn.commit() def recent(self, limit: int = 200) -> list[dict]: """Последние задачи без результатов: для списка нужны только заголовки. Результат каждой задачи весит сотни килобайт, и тянуть их ради списка значит держать в памяти десятки мегабайт ради нескольких строк. """ with self._lock: rows = self._conn.execute( "SELECT id, filename, status, created_at, duration_sec, error " "FROM jobs ORDER BY created_at DESC LIMIT ?", (limit,)).fetchall() return [dict(row) for row in rows] def stats(self) -> dict: with self._lock: rows = self._conn.execute( "SELECT status, COUNT(*) AS n FROM jobs GROUP BY status").fetchall() return {r["status"]: r["n"] for r in rows}