Files
scheduler/backend/watcher.py
T

170 lines
4.8 KiB
Python
Raw Normal View History

2026-06-25 11:31:20 -04:00
from __future__ import annotations
import logging
import os
import threading
import time
from dataclasses import dataclass
from typing import Any
from watchdog.events import FileSystemEvent, FileSystemEventHandler, FileSystemMovedEvent
from watchdog.observers.polling import PollingObserver
2026-07-13 16:17:36 -04:00
from scanner import resolve_runtime_path, scan_results, path_exists_with_smb
2026-06-25 11:31:20 -04:00
LOGGER = logging.getLogger("scheduler.watcher")
@dataclass(frozen=True)
class WatchConfig:
dut_dir: str
ref_dir: str
smb_username: str
smb_password: str
smb_domain: str
class _ResultDirEventHandler(FileSystemEventHandler):
def __init__(self, on_new_directory):
super().__init__()
self._on_new_directory = on_new_directory
def on_created(self, event: FileSystemEvent) -> None:
if event.is_directory:
self._on_new_directory(event.src_path)
def on_moved(self, event: FileSystemMovedEvent) -> None:
if event.is_directory:
self._on_new_directory(event.dest_path)
class ResultDirectoryWatcher:
"""Watch DUT/REF result directories and trigger scans when new result folders appear."""
def __init__(self, poll_interval_seconds: float = 2.0, min_scan_interval_seconds: float = 1.0):
self._poll_interval_seconds = poll_interval_seconds
self._min_scan_interval_seconds = min_scan_interval_seconds
self._lock = threading.RLock()
self._scan_lock = threading.Lock()
self._observer: PollingObserver | None = None
self._current_config: WatchConfig | None = None
self._last_scan_monotonic = 0.0
def configure_from_settings(self, settings: dict[str, Any]) -> None:
dut_dir = _normalize_watch_path(settings.get("dutResultDir"))
ref_dir = _normalize_watch_path(settings.get("refResultDir"))
new_config = WatchConfig(
dut_dir=dut_dir,
ref_dir=ref_dir,
smb_username=str(settings.get("smbUsername") or "").strip(),
smb_password=str(settings.get("smbPassword") or ""),
smb_domain=str(settings.get("smbDomain") or "").strip(),
)
with self._lock:
if new_config == self._current_config and self._observer is not None:
return
self._stop_locked()
self._current_config = new_config
if not dut_dir or not ref_dir:
LOGGER.info("result watcher disabled: DUT/REF result directories are not both set")
return
handler = _ResultDirEventHandler(self._on_new_directory)
observer = PollingObserver(timeout=self._poll_interval_seconds)
2026-07-13 16:17:36 -04:00
smb_creds = {
"username": new_config.smb_username,
"password": new_config.smb_password,
"domain": new_config.smb_domain,
}
2026-06-25 11:31:20 -04:00
scheduled_count = 0
for result_dir, label in ((dut_dir, "DUT"), (ref_dir, "REF")):
2026-07-13 16:17:36 -04:00
if not path_exists_with_smb(result_dir, smb_credentials=smb_creds):
2026-06-25 11:31:20 -04:00
LOGGER.warning("%s result directory does not exist yet, skipping watch: %s", label, result_dir)
continue
observer.schedule(handler, result_dir, recursive=False)
scheduled_count += 1
if scheduled_count == 0:
LOGGER.warning("result watcher not started: no valid result directories to watch")
return
observer.start()
self._observer = observer
LOGGER.info("result watcher started for DUT=%s REF=%s", dut_dir, ref_dir)
self.scan_now(reason="startup")
def _on_new_directory(self, directory_path: str) -> None:
LOGGER.info("new result directory detected: %s", directory_path)
self.scan_now(reason="directory-created")
def scan_now(self, reason: str = "manual") -> None:
with self._lock:
config = self._current_config
if config is None or not config.dut_dir or not config.ref_dir:
return
now = time.monotonic()
if now - self._last_scan_monotonic < self._min_scan_interval_seconds:
return
self._last_scan_monotonic = now
if not self._scan_lock.acquire(blocking=False):
return
def _run_scan() -> None:
credentials = {
"username": config.smb_username,
"password": config.smb_password,
"domain": config.smb_domain,
}
try:
scan_results(config.dut_dir, config.ref_dir, smb_credentials=credentials)
LOGGER.info("result scan finished (%s)", reason)
except Exception:
LOGGER.exception("result scan failed (%s)", reason)
finally:
self._scan_lock.release()
threading.Thread(target=_run_scan, daemon=True).start()
def stop(self) -> None:
with self._lock:
self._stop_locked()
def _stop_locked(self) -> None:
observer = self._observer
self._observer = None
if observer is not None:
observer.stop()
observer.join(timeout=5)
LOGGER.info("result watcher stopped")
_WATCHER = ResultDirectoryWatcher()
def configure_result_watcher(settings: dict[str, Any]) -> None:
_WATCHER.configure_from_settings(settings)
def stop_result_watcher() -> None:
_WATCHER.stop()
def _normalize_watch_path(path_value: Any) -> str:
if path_value is None:
return ""
cleaned = str(path_value).strip()
if not cleaned:
return ""
return str(resolve_runtime_path(cleaned)).strip()
2026-07-13 16:17:36 -04:00