fix: make YAML config persistence atomic and add MJPEG idle keepalive heartbeat

This commit is contained in:
andrew committed 2026-08-19 15:26:28 +07:00
1 parent d9daf0ab01
commit 2205330679
2 files changed
+47 -21

No files matched your search

+4
View File
@@ -17,6 +17,10 @@ All notable changes to the `chicken-counting-sukawarna-det` project are document
- Configured `pip install -e .` in environment setup to resolve `chicken_counter` imports and enable automatic unit test discovery without manual `PYTHONPATH` prefixes. - Configured `pip install -e .` in environment setup to resolve `chicken_counter` imports and enable automatic unit test discovery without manual `PYTHONPATH` prefixes.
### ⚡ Performance & Optimization ### ⚡ Performance & Optimization
- **Atomic Config Persistence (`dashboard.py`)**:
- Added `_config_write_lock` and atomic temporary file replacement (`os.replace`) in `_persist_cycle_start_date` to prevent race conditions or half-written YAML reads.
- **Streaming Socket Keepalive (`dashboard.py`)**:
- Added an idle keepalive heartbeat (`b"--frame\r\n\r\n"`) every 3 seconds to immediately detect and terminate disconnected MJPEG streaming threads during pipeline idle periods.
- **Disk Traversal Caching in Live Dashboard (`dashboard.py`)**: - **Disk Traversal Caching in Live Dashboard (`dashboard.py`)**:
- Replaced repetitive `Path.rglob()` scans on every `/api/mortality/*` endpoint with a thread-safe 15-second in-memory TTL cache (`_MORTALITY_CACHE_TTL = 15.0s`). - Replaced repetitive `Path.rglob()` scans on every `/api/mortality/*` endpoint with a thread-safe 15-second in-memory TTL cache (`_MORTALITY_CACHE_TTL = 15.0s`).
- Added $O(1)$ memory index lookup for serving annotated images (`/api/mortality/image/<filename>`), eliminating full filesystem traversals across 30GB+ video datasets. - Added $O(1)$ memory index lookup for serving annotated images (`/api/mortality/image/<filename>`), eliminating full filesystem traversals across 30GB+ video datasets.
+43 -21
View File
@@ -6,6 +6,7 @@ from __future__ import annotations
import argparse import argparse
import json import json
import mimetypes import mimetypes
import os
import re import re
import sqlite3 import sqlite3
import threading import threading
@@ -27,6 +28,7 @@ TEMPLATE_DIR = Path(__file__).resolve().parent / "templates"
_thread_local = threading.local() _thread_local = threading.local()
_db_lock = threading.Lock() _db_lock = threading.Lock()
_config_write_lock = threading.Lock()
_db_path = "" _db_path = ""
_mortality_dirs: list[Path] = [] _mortality_dirs: list[Path] = []
_mortality_cache_lock = threading.Lock() _mortality_cache_lock = threading.Lock()
@@ -51,29 +53,35 @@ _cycle_start_date: str = _load_initial_cycle_start_date()
def _persist_cycle_start_date(new_date: str) -> bool: def _persist_cycle_start_date(new_date: str) -> bool:
"""Update in-memory cycle_start_date and save to config YAML files.""" """Update in-memory cycle_start_date and save to config YAML files atomically."""
global _cycle_start_date global _cycle_start_date
_cycle_start_date = new_date
updated_any = False updated_any = False
for cfg_name in ("cycle7_batch_optimized.yaml", "cycle7_batch.yaml"): with _config_write_lock:
cfg_path = Path(__file__).resolve().parent / "configs" / cfg_name _cycle_start_date = new_date
if cfg_path.exists(): for cfg_name in ("cycle7_batch_optimized.yaml", "cycle7_batch.yaml"):
content = cfg_path.read_text(encoding="utf-8") cfg_path = Path(__file__).resolve().parent / "configs" / cfg_name
pattern = r"^([ \t]*cycle_start_date:[ \t]*)(?:['\"]?)([^'\"\r\n#]+)(?:['\"]?)([ \t]*(?:#.*)?)$" if cfg_path.exists():
try:
content = cfg_path.read_text(encoding="utf-8")
pattern = r"^([ \t]*cycle_start_date:[ \t]*)(?:['\"]?)([^'\"\r\n#]+)(?:['\"]?)([ \t]*(?:#.*)?)$"
def replacer(match: re.Match) -> str: def replacer(match: re.Match) -> str:
prefix = match.group(1) prefix = match.group(1)
comment = match.group(3) or "" comment = match.group(3) or ""
if comment and not comment.startswith(" "): if comment and not comment.startswith(" "):
comment = f" {comment.lstrip()}" comment = f" {comment.lstrip()}"
if not comment.startswith(" "): if not comment.startswith(" "):
comment = f" {comment}" comment = f" {comment}"
return f'{prefix}"{new_date}"{comment}' return f'{prefix}"{new_date}"{comment}'
new_content, count = re.subn(pattern, replacer, content, count=1, flags=re.MULTILINE) new_content, count = re.subn(pattern, replacer, content, count=1, flags=re.MULTILINE)
if count > 0: if count > 0:
cfg_path.write_text(new_content, encoding="utf-8") tmp_path = cfg_path.with_suffix(f".tmp.{os.getpid()}")
updated_any = True tmp_path.write_text(new_content, encoding="utf-8")
os.replace(tmp_path, cfg_path)
updated_any = True
except Exception as err:
print(f"[dashboard] ⚠️ Error atomically updating {cfg_name}: {err}")
print(f"[dashboard] 📅 Updated cycle_start_date to: {new_date} (persisted in configs: {updated_any})") print(f"[dashboard] 📅 Updated cycle_start_date to: {new_date} (persisted in configs: {updated_any})")
return updated_any return updated_any
@@ -299,9 +307,14 @@ class DashboardHandler(SimpleHTTPRequestHandler):
self.send_header("Cache-Control", "no-cache") self.send_header("Cache-Control", "no-cache")
self.end_headers() self.end_headers()
last_mtime = 0 last_mtime = 0.0
last_write_time = time.time()
keepalive_interval = 3.0
try: try:
while True: while True:
now = time.time()
frame_written = False
try: try:
mtime = frame_path.stat().st_mtime mtime = frame_path.stat().st_mtime
if mtime != last_mtime: if mtime != last_mtime:
@@ -314,10 +327,19 @@ class DashboardHandler(SimpleHTTPRequestHandler):
data + b"\r\n" data + b"\r\n"
) )
self.wfile.flush() self.wfile.flush()
last_write_time = now
frame_written = True
except (FileNotFoundError, OSError): except (FileNotFoundError, OSError):
pass pass
# If idle (no new frames), send a lightweight keepalive ping to detect disconnected clients
if not frame_written and (now - last_write_time) >= keepalive_interval:
self.wfile.write(b"--frame\r\n\r\n")
self.wfile.flush()
last_write_time = now
time.sleep(0.1) time.sleep(0.1)
except (BrokenPipeError, ConnectionResetError): except (BrokenPipeError, ConnectionResetError, OSError):
pass pass
def _handle_shm(self, path): def _handle_shm(self, path):