""" 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() with self.state_lock: self._fill_missing_days() counting_date = self.get_counting_date() if self.current_state is None: self._start_new_day(counting_date) self._persist_day() 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']) # Track IDs are process-local and restart from 1 after every process # start. Persisted dedup keys like "5_in" would silently block new # counts that reuse those IDs, so clear them on resume while keeping # the day's totals. stale_ids = state.get('counted_event_ids') or [] if stale_ids: self.log( f"Cleared {len(stale_ids)} persisted track dedup keys " f"(track IDs reset on restart)" ) state['counted_event_ids'] = [] self.log( f"Resumed {current_date} with total={state['count']} " f"(in={state['count_in']} out={state['count_out']})" ) with open(self.state_file, 'w', encoding='utf-8') as f: json.dump(state, f, indent=2, ensure_ascii=False) 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 _ensure_day_row(self, counting_date, start_time=None, end_time=None): """Insert a zero row for counting_date if it does not already exist.""" now = datetime.now().isoformat() cur = self.db.cursor() cur.execute( """ INSERT OR IGNORE INTO daily_counters (counting_date, camera_name, object_label, total_count, total_in, total_out, start_time, end_time) VALUES (?, ?, ?, 0, 0, 0, ?, ?) """, ( counting_date, self.camera_name, self.object_label, start_time or now, end_time or now, ), ) self.db.commit() def _fill_missing_days(self): """Backfill any missing counting dates from first DB row through today as 0.""" today = self.get_counting_date() cur = self.db.cursor() cur.execute( """ SELECT counting_date FROM daily_counters WHERE camera_name = ? AND object_label = ? ORDER BY counting_date ASC """, (self.camera_name, self.object_label), ) existing = {row[0] for row in cur.fetchall()} if not existing: self._ensure_day_row(today) return start = datetime.strptime(min(existing), '%Y-%m-%d').date() end = datetime.strptime(today, '%Y-%m-%d').date() filled = 0 d = start while d <= end: key = d.isoformat() if key not in existing: self._ensure_day_row(key) filled += 1 d += timedelta(days=1) if filled: self.log( f'Backfilled {filled} zero-activity day(s) ' f'{start.isoformat()}..{end.isoformat()}' ) def record_object_crossing(self, track_id, direction): """Record an object crossing a counting line. direction: 'in' | 'out'. Returns (total_count, day_started, counted). """ 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 counted = False 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) counted = True 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, counted 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: self._fill_missing_days() counting_date = self.get_counting_date() if self.current_state is None: self._start_new_day(counting_date) self._persist_day() continue if self.current_state['counting_date'] != counting_date: self.log('Daily cutoff reached - finalizing day totals') self._persist_day() self.current_state = None self.save_state() self._start_new_day(counting_date) self._persist_day() 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()