from __future__ import annotations import logging import os import shutil import socket import threading import time import uuid from pathlib import Path from app.capture import CaptureService from app.config import Settings from app.db import Database logger = logging.getLogger(__name__) class JobRunner: """SQLite-backed local workers. Every process safely claims jobs through a database lease.""" def __init__(self, database: Database, settings: Settings) -> None: self.database = database self.settings = settings self.stop_event = threading.Event() self.threads: list[threading.Thread] = [] self.maintenance_thread: threading.Thread | None = None self._last_expiry_sweep = 0.0 def start(self) -> None: if not self.maintenance_thread: self.maintenance_thread = threading.Thread( target=self._run_maintenance, name="siteharbor-maintenance", daemon=True, ) self.maintenance_thread.start() if self.settings.worker_enabled and not self.threads: for index in range(self.settings.worker_concurrency): worker_id = f"{socket.gethostname()}-{os.getpid()}-{index}-{uuid.uuid4().hex[:8]}" thread = threading.Thread( target=self._run_worker, args=(worker_id,), name=f"siteharbor-worker-{index}", daemon=True, ) thread.start() self.threads.append(thread) def stop(self) -> None: self.stop_event.set() for thread in self.threads: thread.join(timeout=5) if self.maintenance_thread: self.maintenance_thread.join(timeout=2) def _run_worker(self, worker_id: str) -> None: service = CaptureService(self.database, self.settings, worker_id) while not self.stop_event.is_set(): try: job = self.database.claim_next(worker_id, self.settings.lease_seconds) if not job: self.stop_event.wait(self.settings.queue_poll_seconds) continue service.run(job) except Exception: logger.exception("SiteHarbor worker loop failed") self.stop_event.wait(self.settings.queue_poll_seconds) def _run_maintenance(self) -> None: while not self.stop_event.is_set(): try: self._expire_artifacts_if_due() self._sweep_orphaned_attempts() except Exception: logger.exception("SiteHarbor maintenance loop failed") self.stop_event.wait(10) def _expire_artifacts_if_due(self) -> None: now = time.monotonic() if now - self._last_expiry_sweep < 60: return self._last_expiry_sweep = now for job in self.database.expired_ready_jobs(): self._unlink_under_root(job.get("artifact_path"), self.settings.artifacts_dir) self._unlink_under_root(job.get("report_path"), self.settings.reports_dir) self.database.mark_expired(str(job["id"])) def _sweep_orphaned_attempts(self) -> None: cutoff = time.time() - max(7_200, self.settings.max_duration_seconds_cap + 3_600) referenced_paths = {Path(path).resolve() for path in self.database.referenced_data_paths()} for path in self.settings.work_dir.iterdir(): try: if path.stat().st_mtime >= cutoff: continue if path.is_dir(): shutil.rmtree(path, ignore_errors=True) except OSError: logger.warning("Could not remove orphaned SiteHarbor work directory: %s", path) for root in (self.settings.artifacts_dir, self.settings.reports_dir): for path in root.iterdir(): try: if path.resolve() in referenced_paths or path.stat().st_mtime >= cutoff: continue if path.is_file(): path.unlink(missing_ok=True) except OSError: logger.warning("Could not remove orphaned SiteHarbor output: %s", path) @staticmethod def _unlink_under_root(path_value: object, root: Path) -> None: if not isinstance(path_value, str) or not path_value: return path = Path(path_value) try: if path.resolve().is_relative_to(root.resolve()): path.unlink(missing_ok=True) except OSError: logger.warning("Could not remove expired SiteHarbor artifact: %s", path)