diff --git a/dashboard/src/components/ConfigModal.jsx b/dashboard/src/components/ConfigModal.jsx index 9a3b41e..53ffb50 100644 --- a/dashboard/src/components/ConfigModal.jsx +++ b/dashboard/src/components/ConfigModal.jsx @@ -1,4 +1,5 @@ import { useState, useEffect } from 'react' +import { useQueryClient } from '@tanstack/react-query' import { useConfig, useSaveConfig } from '../hooks/useConfig' import { apiFetch } from '../lib/api' @@ -9,6 +10,7 @@ function fmtSeconds(s) { } export default function ConfigModal({ onClose }) { + const queryClient = useQueryClient() const { data: config, isLoading } = useConfig() const { mutate: save, isPending } = useSaveConfig() @@ -81,6 +83,8 @@ export default function ConfigModal({ onClose }) { avg_time_p3p: form.avg_time_p3p ? String(parseFloat(form.avg_time_p3p) * 60) : null, }), }) + queryClient.invalidateQueries({ queryKey: ['config'] }) + queryClient.invalidateQueries({ queryKey: ['stats'] }) } catch (err) { setSaveError(err?.message ?? 'Failed to save') } finally { @@ -97,6 +101,8 @@ export default function ConfigModal({ onClose }) { if (data?.testCount !== null && data?.testCount !== undefined) { setScanResult({ testCount: data.testCount, completedCount: data.completedCount }) } + queryClient.invalidateQueries({ queryKey: ['stats'] }) + queryClient.invalidateQueries({ queryKey: ['tests'] }) } catch (err) { setSaveError(err?.message ?? 'Rescan failed') } finally { diff --git a/dashboard/src/hooks/useStats.js b/dashboard/src/hooks/useStats.js index e1b6dad..5837378 100644 --- a/dashboard/src/hooks/useStats.js +++ b/dashboard/src/hooks/useStats.js @@ -1,42 +1,14 @@ -import { useQuery, useQueryClient } from '@tanstack/react-query' -import { useEffect } from 'react' -import { getScanStatus, getStats } from '../lib/api' +import { useQuery } from '@tanstack/react-query' +import { getStats } from '../lib/api' export function useStats() { - const queryClient = useQueryClient() - - // Subscribe to SSE updates once; invalidate both stats and tests on any update - useEffect(() => { - const es = new EventSource('/api/events') - es.onmessage = (e) => { - try { - const data = JSON.parse(e.data) - if (data.type === 'update') { - queryClient.invalidateQueries({ queryKey: ['stats'] }) - queryClient.invalidateQueries({ queryKey: ['tests'] }) - } - } catch { /* ignore malformed */ } - } - return () => es.close() - }, [queryClient]) - - const scanStatusQuery = useQuery({ - queryKey: ['scanStatus'], - queryFn: getScanStatus, - refetchInterval: (query) => (query.state.data?.scanning ? 1000 : 5000), - }) - - const isScanning = scanStatusQuery.data?.scanning ?? true - const statsQuery = useQuery({ queryKey: ['stats'], queryFn: getStats, - enabled: !isScanning, - refetchInterval: isScanning ? false : 30_000, }) return { ...statsQuery, - isScanning, + isScanning: false, } } diff --git a/dashboard/src/lib/api.js b/dashboard/src/lib/api.js index 7fa1ada..e69b048 100644 --- a/dashboard/src/lib/api.js +++ b/dashboard/src/lib/api.js @@ -13,7 +13,6 @@ export async function apiFetch(path, options = {}) { } export const getStats = () => apiFetch('/stats') -export const getScanStatus = () => apiFetch('/scan-status') export const getTests = (params = {}) => { const qs = new URLSearchParams( Object.entries(params).filter(([, v]) => v !== '' && v !== undefined && v !== null) diff --git a/server/app.py b/server/app.py index b4a0049..e257163 100644 --- a/server/app.py +++ b/server/app.py @@ -1,14 +1,11 @@ import os -import threading from pathlib import Path -from flask import Flask, Response, jsonify, request, send_from_directory +from flask import Flask, jsonify, request, send_from_directory from flask_cors import CORS from db_py import count_tests, del_config, get_all_tests, get_config, set_config from scanner import full_scan, is_scan_in_progress, resolve_runtime_path, scan_results_only -from sse_py import broadcast, stream_events -from watcher import start_watching PORT = int(os.getenv("PORT", "3001")) ALLOWED_KEYS = { @@ -45,39 +42,6 @@ def _apply_smb_env_from_config(): os.environ[env_key] = str(value) -def _start_full_scan_background(target_dir, results_dir, results_dir_ref, source_label): - def _job(): - try: - full_scan(target_dir, results_dir, results_dir_ref) - start_watching(target_dir, results_dir, results_dir_ref) - tests = get_all_tests() - completed = len([t for t in tests if t.get("completed")]) - print(f"[{source_label}] Scan complete -> {len(tests)} tests, {completed} completed") - except Exception as exc: - print(f"[{source_label}] background fullScan error: {exc}") - finally: - broadcast({"type": "update"}) - - worker = threading.Thread(target=_job, daemon=True) - worker.start() - - -def _start_results_scan_background(results_dir, results_dir_ref, source_label): - def _job(): - try: - scan_results_only(results_dir, results_dir_ref) - tests = get_all_tests() - completed = len([t for t in tests if t.get("completed")]) - print(f"[{source_label}] Results scan complete -> {len(tests)} tests, {completed} completed") - except Exception as exc: - print(f"[{source_label}] background rescan-results error: {exc}") - finally: - broadcast({"type": "update"}) - - worker = threading.Thread(target=_job, daemon=True) - worker.start() - - @app.get("/api/tests") def get_tests_route(): completed = request.args.get("completed") @@ -276,10 +240,11 @@ def set_config_route(): if is_scan_in_progress(): return jsonify({"ok": True, "scanning": True, "testCount": None, "completedCount": None}) - _start_full_scan_background(target_dir, results_dir, results_dir_ref, "config") - return jsonify({"ok": True, "scanning": True, "testCount": None, "completedCount": None}) + full_scan(target_dir, results_dir, results_dir_ref) + tests = get_all_tests() + completed = len([t for t in tests if t.get("completed")]) + return jsonify({"ok": True, "scanning": False, "testCount": len(tests), "completedCount": completed}) - broadcast({"type": "update"}) return jsonify({"ok": True, "testCount": None, "completedCount": None}) @@ -295,8 +260,10 @@ def rescan_route(): if is_scan_in_progress(): return jsonify({"ok": True, "scanning": True, "testCount": None, "completedCount": None}) - _start_full_scan_background(target_dir, results_dir, results_dir_ref, "config") - return jsonify({"ok": True, "scanning": True, "testCount": None, "completedCount": None}) + full_scan(target_dir, results_dir, results_dir_ref) + tests = get_all_tests() + completed = len([t for t in tests if t.get("completed")]) + return jsonify({"ok": True, "scanning": False, "testCount": len(tests), "completedCount": completed}) @app.post("/api/config/rescan-results") @@ -310,18 +277,10 @@ def rescan_results_route(): if is_scan_in_progress(): return jsonify({"ok": True, "scanning": True, "testCount": None, "completedCount": None}) - _start_results_scan_background(results_dir, results_dir_ref, "config") - return jsonify({"ok": True, "scanning": True, "testCount": None, "completedCount": None}) - - -@app.get("/api/events") -def events_route(): - headers = { - "Cache-Control": "no-cache", - "Connection": "keep-alive", - "X-Accel-Buffering": "no", - } - return Response(stream_events(), mimetype="text/event-stream", headers=headers) + scan_results_only(results_dir, results_dir_ref) + tests = get_all_tests() + completed = len([t for t in tests if t.get("completed")]) + return jsonify({"ok": True, "scanning": False, "testCount": len(tests), "completedCount": completed}) @app.get("/") @@ -351,9 +310,6 @@ def bootstrap(): tests = get_all_tests() completed = len([t for t in tests if t.get("completed")]) print(f"[server] Scan complete -> {len(tests)} tests found, {completed} completed") - - start_watching(target_dir, results_dir) - print("[server] Watching for changes.") else: print("[server] No directories configured -> open the dashboard settings to get started.") diff --git a/server/dashboard.db b/server/dashboard.db index 86b8ee1..bff3e1c 100644 Binary files a/server/dashboard.db and b/server/dashboard.db differ diff --git a/server/db_py.py b/server/db_py.py index 765d167..d956262 100644 --- a/server/db_py.py +++ b/server/db_py.py @@ -248,6 +248,19 @@ def reset_by_file_id_and_device(test_id, device): ) +def reset_all_results_state(): + with _tx(): + _conn.execute( + """ + UPDATE tests + SET completed = 0, + completed_at = NULL, + duration_seconds = NULL, + tput_results = NULL + """ + ) + + def clear_tests(): with _tx(): _conn.execute("DELETE FROM tests") diff --git a/server/requirements.txt b/server/requirements.txt index 3be0ebb..bad148b 100644 --- a/server/requirements.txt +++ b/server/requirements.txt @@ -1,4 +1,3 @@ Flask>=3.0.0,<4.0.0 Flask-Cors>=4.0.1,<5.0.0 -watchdog>=4.0.1,<5.0.0 smbprotocol>=1.13.0,<2.0.0 diff --git a/server/scanner.py b/server/scanner.py index f71639d..84cbf36 100644 --- a/server/scanner.py +++ b/server/scanner.py @@ -7,6 +7,7 @@ from db_py import ( clear_tests, get_config, mark_completed, + reset_all_results_state, update_all_p2p_coe_pairs_sql, update_all_p3p_pairs_sql, upsert_test, @@ -289,6 +290,9 @@ def scan_results_only(results_dir, results_dir_ref): print(f"[scanner] results-only scan dir : {results_dir}") print(f"[scanner] results-only scan dir ref : {results_dir_ref}") + # Rebuild completion state from current result folders so deletions are reflected. + reset_all_results_state() + scan_results(results_dir) scan_results(results_dir_ref) coe_pairs = update_all_p2p_coe_pairs_sql() diff --git a/server/sse_py.py b/server/sse_py.py deleted file mode 100644 index 02b7681..0000000 --- a/server/sse_py.py +++ /dev/null @@ -1,32 +0,0 @@ -import json -import queue -import threading - -_clients = set() -_clients_lock = threading.Lock() - - -def stream_events(): - q = queue.Queue() - with _clients_lock: - _clients.add(q) - - try: - yield 'data: {"type":"connected"}\n\n' - while True: - try: - payload = q.get(timeout=20) - yield f"data: {payload}\n\n" - except queue.Empty: - yield ": heartbeat\n\n" - finally: - with _clients_lock: - _clients.discard(q) - - -def broadcast(data): - payload = json.dumps(data) - with _clients_lock: - clients = list(_clients) - for q in clients: - q.put_nowait(payload) diff --git a/server/watcher.py b/server/watcher.py deleted file mode 100644 index 839d06a..0000000 --- a/server/watcher.py +++ /dev/null @@ -1,270 +0,0 @@ -import os -import threading - -from watchdog.events import FileSystemEventHandler -from watchdog.observers import Observer - -from db_py import reset_by_file_id_and_device -from scanner import ( - full_scan, - parse_deleted_result_dir_name, - process_result_dir, - _is_unc_path, - _normalize_input_path, -) -from sse_py import broadcast - -_target_observer = None -_results_observers = [] -_scan_timer = None -_scan_lock = threading.Lock() - - -def _path_exists(path): - """Check if path exists, handling UNC paths.""" - if _is_unc_path(path): - try: - import smbclient - return smbclient.path.isdir(path) - except Exception: - return False - return os.path.isdir(path) - - -def _normalize_watcher_path(path): - """Normalize path for watcher comparison, handling UNC paths.""" - normalized = _normalize_input_path(path) - if not _is_unc_path(normalized): - normalized = os.path.normcase(os.path.abspath(normalized)) - return normalized - - -def _schedule_full_scan(target_dir, results_dir, delay_seconds=1.0): - global _scan_timer - with _scan_lock: - if _scan_timer: - _scan_timer.cancel() - - def run_scan(): - try: - full_scan(target_dir, results_dir) - broadcast({"type": "update"}) - except Exception as exc: - print(f"[watcher] fullScan error: {exc}") - - _scan_timer = threading.Timer(delay_seconds, run_scan) - _scan_timer.daemon = True - _scan_timer.start() - - -class _TargetHandler(FileSystemEventHandler): - def __init__(self, target_dir, results_dir): - self.target_dir = os.path.normcase(os.path.abspath(target_dir)) - self.results_dir = results_dir - - def _is_target_test_file(self, path): - normalized = os.path.normcase(os.path.abspath(path)) - relative = os.path.relpath(normalized, self.target_dir) - parts = relative.split(os.sep) - if len(parts) != 2: - return False - - filename = parts[-1] - return filename.endswith(".ini") and not filename.startswith("GLOBAL") - - def on_created(self, event): - if not event.is_directory and self._is_target_test_file(event.src_path): - _schedule_full_scan(self.target_dir, self.results_dir) - - def on_deleted(self, event): - if not event.is_directory and self._is_target_test_file(event.src_path): - _schedule_full_scan(self.target_dir, self.results_dir) - - def on_moved(self, event): - if event.is_directory: - return - - # On Windows, create/delete in Explorer can show up as move/rename events. - if self._is_target_test_file(event.src_path) or self._is_target_test_file(event.dest_path): - _schedule_full_scan(self.target_dir, self.results_dir) - - -class _ResultsHandler(FileSystemEventHandler): - def __init__(self, results_dir): - self.results_dir = _normalize_watcher_path(results_dir) - - def _is_direct_child_dir(self, path): - try: - if _is_unc_path(path): - parent = _normalize_watcher_path(os.path.dirname(path)) - else: - parent = os.path.normcase(os.path.abspath(os.path.dirname(path))) - return parent == self.results_dir - except Exception: - return False - - def _is_file_under_result_child(self, path): - try: - if _is_unc_path(path): - parent_dir = _normalize_watcher_path(os.path.dirname(path)) - grandparent = _normalize_watcher_path(os.path.dirname(parent_dir)) - else: - parent_dir = os.path.normcase(os.path.abspath(os.path.dirname(path))) - grandparent = os.path.normcase(os.path.abspath(os.path.dirname(parent_dir))) - return grandparent == self.results_dir - except Exception: - return False - - def _result_child_name_for_file(self, path): - return os.path.basename(os.path.dirname(path)) - - def _result_child_name_for_dir(self, path): - return os.path.basename(path) - - def on_created(self, event): - # On Windows, is_directory may be False even for directories (timing issue), - # so check _is_direct_child_dir regardless of the flag. - if self._is_direct_child_dir(event.src_path): - dir_name = self._result_child_name_for_dir(event.src_path) - process_result_dir(self.results_dir, dir_name) - broadcast({"type": "update"}) - return - - if self._is_file_under_result_child(event.src_path): - dir_name = self._result_child_name_for_file(event.src_path) - process_result_dir(self.results_dir, dir_name) - broadcast({"type": "update"}) - - def on_deleted(self, event): - # On Windows, when a directory is deleted watchdog may report is_directory=False - # because os.path.isdir() returns False by the time the event is processed. - # Check _is_direct_child_dir first regardless of the is_directory flag. - if self._is_direct_child_dir(event.src_path): - dir_name = self._result_child_name_for_dir(event.src_path) - test_id, device = parse_deleted_result_dir_name(dir_name) - if test_id: - reset_by_file_id_and_device(test_id, device) - broadcast({"type": "update"}) - return - - if self._is_file_under_result_child(event.src_path): - dir_name = self._result_child_name_for_file(event.src_path) - process_result_dir(self.results_dir, dir_name) - broadcast({"type": "update"}) - - def on_moved(self, event): - src_in_root = self._is_direct_child_dir(event.src_path) - dst_in_root = self._is_direct_child_dir(event.dest_path) - - # Handle directory-level moves regardless of is_directory flag (Windows timing issue). - if src_in_root or dst_in_root: - # Result directory moved out (includes Recycle Bin delete on Windows). - if src_in_root and not dst_in_root: - old_name = self._result_child_name_for_dir(event.src_path) - test_id, device = parse_deleted_result_dir_name(old_name) - if test_id: - reset_by_file_id_and_device(test_id, device) - broadcast({"type": "update"}) - return - - # Result directory moved in. - if dst_in_root and not src_in_root: - new_name = self._result_child_name_for_dir(event.dest_path) - process_result_dir(self.results_dir, new_name) - broadcast({"type": "update"}) - return - - # Result directory renamed within root. - if src_in_root and dst_in_root: - old_name = self._result_child_name_for_dir(event.src_path) - new_name = self._result_child_name_for_dir(event.dest_path) - test_id, device = parse_deleted_result_dir_name(old_name) - if test_id: - reset_by_file_id_and_device(test_id, device) - process_result_dir(self.results_dir, new_name) - broadcast({"type": "update"}) - return - - src_file_in_result = self._is_file_under_result_child(event.src_path) - dst_file_in_result = self._is_file_under_result_child(event.dest_path) - - if src_file_in_result: - src_name = self._result_child_name_for_file(event.src_path) - process_result_dir(self.results_dir, src_name) - - if dst_file_in_result: - dst_name = self._result_child_name_for_file(event.dest_path) - if not src_file_in_result or src_name != dst_name: - process_result_dir(self.results_dir, dst_name) - - if src_file_in_result or dst_file_in_result: - broadcast({"type": "update"}) - - -def stop_watching(): - global _target_observer, _results_observers, _scan_timer - - if _scan_timer: - _scan_timer.cancel() - _scan_timer = None - - if _target_observer: - _target_observer.stop() - _target_observer.join(timeout=2) - _target_observer = None - - for obs in _results_observers: - try: - obs.stop() - obs.join(timeout=2) - except Exception as e: - print(f"[watcher] error stopping results observer: {e}") - _results_observers.clear() - - -def start_watching(target_dir, results_dir, results_dir_ref=None): - global _target_observer, _results_observers - - stop_watching() - if not target_dir or not results_dir: - print("[watcher] target_dir/results_dir not configured; watcher disabled.") - return - - if not _path_exists(target_dir): - print(f"[watcher] target_dir does not exist or is not accessible: {target_dir}") - print("[watcher] watcher disabled until valid paths are configured.") - return - - if not _path_exists(results_dir): - print(f"[watcher] results_dir does not exist or is not accessible: {results_dir}") - print("[watcher] watcher disabled until valid paths are configured.") - return - - try: - _target_observer = Observer() - _target_observer.schedule(_TargetHandler(target_dir, results_dir), target_dir, recursive=True) - _target_observer.daemon = True - _target_observer.start() - print(f"[watcher] target observer started for: {target_dir}") - - # Watch primary results directory - results_obs = Observer() - results_obs.schedule(_ResultsHandler(results_dir), results_dir, recursive=True) - results_obs.daemon = True - results_obs.start() - _results_observers.append(results_obs) - print(f"[watcher] results observer started for: {results_dir}") - - # Watch reference results directory if provided and different - if results_dir_ref and results_dir_ref != results_dir and _path_exists(results_dir_ref): - ref_obs = Observer() - ref_obs.schedule(_ResultsHandler(results_dir_ref), results_dir_ref, recursive=True) - ref_obs.daemon = True - ref_obs.start() - _results_observers.append(ref_obs) - print(f"[watcher] results observer started for reference: {results_dir_ref}") - elif results_dir_ref and results_dir_ref != results_dir: - print(f"[watcher] reference results_dir does not exist or is not accessible: {results_dir_ref}") - except Exception as exc: - print(f"[watcher] failed to start watchers: {exc}") - stop_watching()