Files
2026-09-05 06:07:25 +05:30

122 lines
4.6 KiB
Python

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)