mydocker/side-agent/snapit.py

190 lines
7.2 KiB
Python

#!/usr/bin/env python3
import os
import filecmp
import time
import shutil
import threading
import glob
from datetime import datetime
import requests
import logging
# Konfiguracja logowania
logging.basicConfig(level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)
# Konfiguracja zmiennych (z opcją nadpisania przez ENV)
TIMEOUT = int(os.getenv("SNAPIT_TIMEOUT", 3))
EXTENSION = "jpg"
CAMERAS = ["reolink_1", "reolink_2", "reolink_3", "reolink_4", "hikvision_2", "reolink_5", "hikvision_4"]
OBJECTS = ["person", "car", "dog", "cat", "animal"]
SOURCE_DIR = os.getenv("SOURCE_DIR", "/source")
TARGET_DIR = os.getenv("TARGET_DIR", "/target")
NOTIFICATION_DIR = os.path.join(TARGET_DIR, "frigate-detections")
PUSHGATEWAY_HOST = os.getenv("PUSHGATEWAY_HOST", "192.168.1.132")
PUSHGATEWAY_PORT = os.getenv("PUSHGATEWAY_PORT", "9091")
# URL bazowy - kluczem jest to, że każda kamera będzie miała swój unikalny "instance" lub label w URL
BASE_PUSH_URL = f"http://{PUSHGATEWAY_HOST}:{PUSHGATEWAY_PORT}/metrics/job/camera_detection"
COUNTER = 0
DETECTION_COUNTS = {camera: 0 for camera in CAMERAS}
DETECTION_TIMERS = {}
# Blokada dla operacji na słowniku DETECTION_COUNTS (wątkowość)
COUNT_LOCK = threading.Lock()
# Upewnij się, że katalog powiadomień istnieje
os.makedirs(NOTIFICATION_DIR, exist_ok=True)
logger.info(f"snapit.py started. Pushgateway: {BASE_PUSH_URL}")
def push_metric(camera, count):
"""
Push metric to Prometheus Pushgateway.
UWAGA: Dodajemy '/camera/{camera}' do URL, aby każda kamera była traktowana
jako osobna grupa metryk. Inaczej kamery nadpisywałyby się nawzajem.
"""
try:
# Używamy labela 'camera' jako klucza grupowania w URL
url = f"{BASE_PUSH_URL}/camera/{camera}"
metric_data = f"# TYPE person_detected gauge\nperson_detected {{camera=\"{camera}\"}} {count}\n"
response = requests.post(url, data=metric_data, timeout=5)
if response.status_code != 200:
logger.error(f"Failed to push metric for {camera}. Status: {response.status_code}, Resp: {response.text}")
else:
logger.debug(f"Pushed metric for {camera}: {count}")
except Exception as e:
logger.error(f"Error pushing metric for {camera}: {e}")
def change_detection_count(camera, amount):
"""Thread-safe update of detection count."""
with COUNT_LOCK:
DETECTION_COUNTS[camera] += amount
# Zabezpieczenie przed ujemnymi wartościami (choć logika na to nie powinna pozwolić)
if DETECTION_COUNTS[camera] < 0:
DETECTION_COUNTS[camera] = 0
current_count = DETECTION_COUNTS[camera]
action = "Increased" if amount > 0 else "Decreased"
logger.info(f"{action} detection count for {camera}. Current count: {current_count}")
push_metric(camera, current_count)
return current_count
def increase_detection(camera):
"""Increase detection count and reset the timer."""
change_detection_count(camera, 1)
# Reset timera (anuluj stary, uruchom nowy)
# Lock nie jest tu konieczny dla timera, bo to operacja atomowa na referencji,
# ale dla porządku w logice wątków warto uważać.
if camera in DETECTION_TIMERS:
try:
DETECTION_TIMERS[camera].cancel()
except Exception:
pass # Ignoruj jeśli timer już wygasł
# Uruchom nowy timer, który zmniejszy licznik po 30s
timer = threading.Timer(30, decrease_detection, [camera])
DETECTION_TIMERS[camera] = timer
timer.start()
def decrease_detection(camera):
"""Decrease detection count after timeout."""
# Ta funkcja jest wywoływana z osobnego wątku!
change_detection_count(camera, -1)
def process_camera(camera):
"""Check for new files in the source directory and handle detection."""
pattern = os.path.join(SOURCE_DIR, f"*{camera}*.{EXTENSION}")
try:
# Pobieramy pliki. Optymalizacja: jeśli plików jest masa, to może być wolne.
# W przyszłości można rozważyć `os.scandir`.
matching_files = glob.glob(pattern)
if not matching_files:
return
# Sortowanie po dacie modyfikacji (najnowsze pierwsze)
# Optymalizacja: pobieramy tylko najnowszy plik bez sortowania całej listy jeśli to możliwe,
# ale max(key=...) jest szybsze niż sort() jeśli potrzebujemy tylko jednego.
newest_file = max(matching_files, key=os.path.getmtime)
# Ignoruj pliki 'clean' (jeśli logika tego wymaga)
if 'clean' in newest_file:
return
file_found = newest_file
# Rozpoznawanie obiektu
file_basename = os.path.basename(file_found)
detected_object = 'person' # Domyślnie
# Szybsze sprawdzanie substringów
detected_object_found = next((obj for obj in OBJECTS if obj in file_basename.lower()), 'person')
# Ścieżka pliku powiadomienia dla HA
notification_file = os.path.join(NOTIFICATION_DIR, f"{camera}-{detected_object_found}.{EXTENSION}")
# Sprawdzenie czy plik jest nowy/zmieniony
is_changed = False
if not os.path.isfile(notification_file):
is_changed = True
elif not filecmp.cmp(file_found, notification_file, shallow=True):
# shallow=True sprawdza rozmiar i czas, co jest znacznie szybsze niż czytanie treści (filecmp domyślnie czyta)
# Jeśli jednak pliki mają identyczny rozmiar/czas ale inną treść (mało prawdopodobne przy kamerach), daj shallow=False
is_changed = True
if is_changed:
logger.info(f"New detection for {camera}: {detected_object_found}")
# Kopiowanie do pliku notyfikacji (nadpisanie)
# shutil.copy2 kopiuje też metadane (czas), co jest ważne dla filecmp
shutil.copy2(file_found, notification_file)
# Archiwizacja
date_str = datetime.now().strftime('%Y-%m-%d')
time_str = datetime.now().strftime('%Y%m%d-%H%M%S')
# Używamy os.path.join dla kompatybilności
archive_dir = f"/hikvision/objects-detected/{camera}/{date_str}"
os.makedirs(archive_dir, exist_ok=True)
archive_file = f"{archive_dir}/detected_on_{camera}_{detected_object_found}_{time_str}.{EXTENSION}"
shutil.copy2(file_found, archive_file)
# Zwiększ licznik
increase_detection(camera)
except Exception as e:
logger.error(f"Error processing camera {camera}: {str(e)}")
def main_loop():
global COUNTER
logger.info("Starting main loop...")
while True:
try:
for camera in CAMERAS:
process_camera(camera)
COUNTER += 1
# Log keepalive co 300 cykli (ok. 15 minut przy timeout 3s)
if COUNTER % 300 == 0:
logger.info(f"Keepalive check. Loop count: {COUNTER}")
except Exception as e:
logger.error(f"Critical error in main loop: {e}")
time.sleep(TIMEOUT)
if __name__ == "__main__":
main_loop()