""" Production daily counter persistence for edge counter. No batches: tracks counter_in / counter_out and the daily total per counting day, delimited by the daily cutoff time. SQLite schema + current_counter.json state. """ import json import sqlite3 import threading import time from datetime import datetime, timedelta from pathlib import Path class CounterStore: def __init__( self, db_path, state_file, camera_name, object_label='object', cutoff_time='20:00', carry_ids=50, logger=print, ): self.db_path = db_path self.state_file = Path(state_file) self.camera_name = camera_name self.object_label = object_label self.cutoff_time_str = cutoff_time datetime.strptime(cutoff_time, '%H:%M') self.carry_ids = int(carry_ids) self.log = logger self.state_lock = threading.Lock() self.shutdown_event = threading.Event() Path(db_path).parent.mkdir(parents=True, exist_ok=True) self.state_file.parent.mkdir(parents=True, exist_ok=True) self.db = sqlite3.connect(db_path, check_same_thread=False) self._init_db() self.current_state = self._load_state() def _init_db(self): cur = self.db.cursor() cur.execute( """ CREATE TABLE IF NOT EXISTS daily_counters ( id INTEGER PRIMARY KEY AUTOINCREMENT, counting_date TEXT NOT NULL, camera_name TEXT NOT NULL, object_label TEXT NOT NULL, total_count INTEGER NOT NULL DEFAULT 0, total_in INTEGER NOT NULL DEFAULT 0, total_out INTEGER NOT NULL DEFAULT 0, start_time TEXT, end_time TEXT, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, UNIQUE(counting_date, camera_name, object_label) ) """ ) self.db.commit() def get_counting_date(self, dt=None): if dt is None: dt = datetime.now() cutoff = datetime.strptime(self.cutoff_time_str, '%H:%M').time() if dt.time() < cutoff: return dt.date().isoformat() return (dt.date() + timedelta(days=1)).isoformat() def _load_state(self): if not self.state_file.exists(): return None try: with open(self.state_file, 'r', encoding='utf-8') as f: state = json.load(f) current_date = self.get_counting_date() if state.get('counting_date') != current_date: self.log( f"State file belongs to previous counting day " f"({state.get('counting_date')}). Starting fresh." ) self.state_file.unlink(missing_ok=True) return None state.setdefault('count_in', 0) state.setdefault('count_out', 0) state.setdefault('count', state['count_in'] + state['count_out']) state.setdefault('counted_event_ids', []) self.log( f"Resumed {current_date} with total={state['count']} " f"(in={state['count_in']} out={state['count_out']})" ) return state except Exception as exc: self.log(f'Failed to load state file: {exc}') return None def save_state(self): if self.current_state is None: self.state_file.unlink(missing_ok=True) return with open(self.state_file, 'w', encoding='utf-8') as f: json.dump(self.current_state, f, indent=2, ensure_ascii=False) def _start_new_day(self, counting_date): now = datetime.now().isoformat() carried = [] if self.current_state is not None: try: carried = self.current_state['counted_event_ids'][-self.carry_ids:] except (KeyError, TypeError): carried = [] self.current_state = { 'counting_date': counting_date, 'count': 0, 'count_in': 0, 'count_out': 0, 'start_time': now, 'last_detection_time': now, 'counted_event_ids': carried, } self.save_state() self.log(f'Started counting day {counting_date} ({self.object_label})') def record_object_crossing(self, track_id, direction): """Record an object crossing a counting line. direction: 'in' | 'out'.""" with self.state_lock: counting_date = self.get_counting_date() day_started = False if self.current_state is None or self.current_state['counting_date'] != counting_date: self._start_new_day(counting_date) day_started = True event_key = f"{track_id}_{direction}" if event_key not in self.current_state['counted_event_ids']: self.current_state['count'] += 1 if direction == 'in': self.current_state['count_in'] += 1 else: self.current_state['count_out'] += 1 self.current_state['counted_event_ids'].append(event_key) self.log( f'Counted {direction} (track {track_id}) | {counting_date} ' f'total: {self.current_state["count"]} ' f'(in={self.current_state["count_in"]} out={self.current_state["count_out"]})' ) self._persist_day() self.current_state['last_detection_time'] = datetime.now().isoformat() self.save_state() return self.current_state['count'], day_started def _persist_day(self): state = self.current_state cur = self.db.cursor() cur.execute( """ INSERT INTO daily_counters (counting_date, camera_name, object_label, total_count, total_in, total_out, start_time, end_time) VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(counting_date, camera_name, object_label) DO UPDATE SET total_count = excluded.total_count, total_in = excluded.total_in, total_out = excluded.total_out, end_time = excluded.end_time, updated_at = CURRENT_TIMESTAMP """, ( state['counting_date'], self.camera_name, self.object_label, state['count'], state['count_in'], state['count_out'], state['start_time'], datetime.now().isoformat(), ), ) self.db.commit() def cutoff_watcher_loop(self): while not self.shutdown_event.is_set(): time.sleep(60) with self.state_lock: if self.current_state is None: continue if self.current_state['counting_date'] != self.get_counting_date(): self.log('Daily cutoff reached - finalizing day totals') self._persist_day() self.current_state = None self.save_state() def start_cutoff_watcher(self): t = threading.Thread(target=self.cutoff_watcher_loop, daemon=True) t.start() return t @property def current_count(self): if self.current_state is None: return 0 return self.current_state['count'] @property def current_count_in(self): if self.current_state is None: return 0 return self.current_state.get('count_in', 0) @property def current_count_out(self): if self.current_state is None: return 0 return self.current_state.get('count_out', 0) def _day_totals(self, counting_date=None): if counting_date is None: counting_date = self.get_counting_date() cur = self.db.cursor() cur.execute( """ SELECT COALESCE(total_count, 0), COALESCE(total_in, 0), COALESCE(total_out, 0) FROM daily_counters WHERE counting_date = ? AND camera_name = ? AND object_label = ? """, (counting_date, self.camera_name, self.object_label), ) row = cur.fetchone() return row if row else (0, 0, 0) def display_total(self): return self._day_totals()[0] def display_in(self): return self._day_totals()[1] def display_out(self): return self._day_totals()[2] def shutdown(self): self.shutdown_event.set() with self.state_lock: if self.current_state is not None: self._persist_day() self.db.close()