190 lines
7.2 KiB
Python
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()
|
|
|