diff --git a/counter_dashboard.py b/counter_dashboard.py index 9cd60b6..753de4a 100644 --- a/counter_dashboard.py +++ b/counter_dashboard.py @@ -77,7 +77,7 @@ def api_current_batch(): "batch_number": None, "counting_date": None, } - ), 404 + ), 200 except Exception as e: return jsonify( { @@ -90,6 +90,42 @@ def api_current_batch(): ), 500 +@app.route("/api/previous-batch") +def api_previous_batch(): + """Get the most recently completed batch.""" + conn = get_db() + cur = conn.cursor() + + cur.execute( + """ + SELECT counting_date, batch_number, count, start_time, end_time, + ROUND( + (julianday(end_time) - julianday(start_time)) * 24 * 60, 1 + ) as duration_minutes + FROM batches + ORDER BY end_time DESC + LIMIT 1 + """ + ) + + row = cur.fetchone() + conn.close() + + if row: + return jsonify( + { + "success": True, + "date": row["counting_date"], + "batch_number": row["batch_number"], + "count": row["count"], + "start_time": row["start_time"], + "end_time": row["end_time"], + "duration_minutes": row["duration_minutes"], + } + ) + return jsonify({"success": False, "error": "No previous batch"}), 200 + + @app.route("/api/summary") def api_summary(): """Get summary statistics for the dashboard cards.""" diff --git a/counter_dashboard.py.20260605 b/counter_dashboard.py.20260605 new file mode 100644 index 0000000..9cd60b6 --- /dev/null +++ b/counter_dashboard.py.20260605 @@ -0,0 +1,462 @@ +#!/usr/bin/env python3 +""" +Frigate Counter Dashboard +A beautiful, interactive Flask web app for viewing daily batch counting data. +""" + +import json +import os +import sqlite3 +import csv +from io import StringIO +from datetime import datetime, timedelta + +from flask import Flask, render_template, jsonify, request, Response +from werkzeug.serving import WSGIRequestHandler + + +app = Flask(__name__, template_folder="templates") +app.config["SECRET_KEY"] = os.getenv("SECRET_KEY", "frigate-counter-dashboard") + +DB_PATH = os.getenv("DB_PATH", "frigate_counter.db") +CURRENT_BATCH_PATH = os.getenv("CURRENT_BATCH_PATH", "current_batch.json") +CUTOFF_TIME = os.getenv("CUTOFF_TIME", "17:00") + + +# ------------------------------------------------------------------ # +# Database helpers +# ------------------------------------------------------------------ # +def get_db(): + conn = sqlite3.connect(DB_PATH) + conn.row_factory = sqlite3.Row + return conn + + +# def get_counting_date(dt=None, cutoff_str="17:00"): +def get_counting_date(dt=None, cutoff_str=CUTOFF_TIME): + """Replicate the service logic for determining the counting date.""" + if dt is None: + dt = datetime.now() + cutoff = datetime.strptime(cutoff_str, "%H:%M").time() + if dt.time() < cutoff: + return dt.date().isoformat() + return (dt.date() + timedelta(days=1)).isoformat() + + +# ------------------------------------------------------------------ # +# Routes +# ------------------------------------------------------------------ # +@app.route("/") +def index(): + """Main dashboard page.""" + return render_template("dashboard.html") + + +@app.route("/api/current-batch") +def api_current_batch(): + """Get real-time data from the current active batch file.""" + try: + with open(CURRENT_BATCH_PATH, "r") as f: + data = json.load(f) + return jsonify( + { + "success": True, + "counting_date": data.get("counting_date"), + "batch_number": data.get("batch_number"), + "count": data.get("count", 0), + "start_time": data.get("start_time"), + "last_detection_time": data.get("last_detection_time"), + } + ) + except FileNotFoundError: + return jsonify( + { + "success": False, + "error": "No active batch", + "count": 0, + "batch_number": None, + "counting_date": None, + } + ), 404 + except Exception as e: + return jsonify( + { + "success": False, + "error": str(e), + "count": 0, + "batch_number": None, + "counting_date": None, + } + ), 500 + + +@app.route("/api/summary") +def api_summary(): + """Get summary statistics for the dashboard cards.""" + conn = get_db() + cur = conn.cursor() + + # Today's counting date (based on cutoff) + today = get_counting_date() + + # Today's stats + cur.execute( + """ + SELECT COALESCE(total_count, 0) as total_count, + COALESCE(total_batches, 0) as total_batches + FROM daily_summaries + WHERE counting_date = ? + """, + (today,), + ) + today_row = cur.fetchone() + + # Yesterday's stats + yesterday = ( + datetime.strptime(today, "%Y-%m-%d").date() - timedelta(days=1) + ).isoformat() + cur.execute( + """ + SELECT COALESCE(total_count, 0) as total_count, + COALESCE(total_batches, 0) as total_batches + FROM daily_summaries + WHERE counting_date = ? + """, + (yesterday,), + ) + yesterday_row = cur.fetchone() + + # All-time totals + cur.execute(""" + SELECT COALESCE(SUM(total_count), 0) as grand_total, + COALESCE(SUM(total_batches), 0) as grand_batches, + COUNT(DISTINCT counting_date) as total_days + FROM daily_summaries + """) + all_time = cur.fetchone() + + # Average per day + cur.execute(""" + SELECT ROUND(AVG(total_count), 1) as avg_per_day + FROM daily_summaries + """) + avg = cur.fetchone() + + # Best day + cur.execute(""" + SELECT counting_date, total_count + FROM daily_summaries + ORDER BY total_count DESC + LIMIT 1 + """) + best = cur.fetchone() + + conn.close() + + return jsonify( + { + "today": { + "date": today, + "total_count": today_row["total_count"] if today_row else 0, + "total_batches": today_row["total_batches"] if today_row else 0, + }, + "yesterday": { + "date": yesterday, + "total_count": yesterday_row["total_count"] if yesterday_row else 0, + "total_batches": yesterday_row["total_batches"] if yesterday_row else 0, + }, + "all_time": { + "grand_total": all_time["grand_total"], + "grand_batches": all_time["grand_batches"], + "total_days": all_time["total_days"], + }, + "average_per_day": avg["avg_per_day"] or 0, + "best_day": { + "date": best["counting_date"] if best else None, + "count": best["total_count"] if best else 0, + }, + } + ) + + +@app.route("/api/daily-data") +def api_daily_data(): + """Get daily data for charts and table.""" + days = request.args.get("days", 30, type=int) + date_from = (datetime.now() - timedelta(days=days)).date().isoformat() + + conn = get_db() + cur = conn.cursor() + + # Daily summaries for chart + cur.execute( + """ + SELECT counting_date, total_count, total_batches, + ROUND(CAST(total_count AS FLOAT) / total_batches, 1) as avg_per_batch + FROM daily_summaries + WHERE counting_date >= ? + ORDER BY counting_date ASC + """, + (date_from,), + ) + + daily_data = [] + for row in cur.fetchall(): + daily_data.append( + { + "date": row["counting_date"], + "total_count": row["total_count"], + "total_batches": row["total_batches"], + "avg_per_batch": row["avg_per_batch"] or 0, + } + ) + + conn.close() + return jsonify(daily_data) + + +@app.route("/api/day-detail/") +def api_day_detail(date): + """Get detailed batch information for a specific day.""" + conn = get_db() + cur = conn.cursor() + + # Batches for this day + cur.execute( + """ + SELECT batch_number, count, start_time, end_time, + ROUND( + (julianday(end_time) - julianday(start_time)) * 24 * 60, 1 + ) as duration_minutes + FROM batches + WHERE counting_date = ? + ORDER BY batch_number ASC + """, + (date,), + ) + + batches = [] + total_duration = 0 + for row in cur.fetchall(): + duration = row["duration_minutes"] or 0 + total_duration += duration + batches.append( + { + "batch_number": row["batch_number"], + "count": row["count"], + "start_time": row["start_time"], + "end_time": row["end_time"], + "duration_minutes": duration, + } + ) + + # Summary for the day + cur.execute( + """ + SELECT total_count, total_batches + FROM daily_summaries + WHERE counting_date = ? + """, + (date,), + ) + summary = cur.fetchone() + + conn.close() + + return jsonify( + { + "date": date, + "total_count": summary["total_count"] if summary else 0, + "total_batches": summary["total_batches"] if summary else 0, + "total_duration_minutes": round(total_duration, 1), + "avg_duration_minutes": round(total_duration / len(batches), 1) + if batches + else 0, + "batches": batches, + } + ) + + +@app.route("/api/recent-batches") +def api_recent_batches(): + """Get the most recent batches across all days.""" + limit = request.args.get("limit", 10, type=int) + + conn = get_db() + cur = conn.cursor() + + cur.execute( + """ + SELECT counting_date, batch_number, count, start_time, end_time, + ROUND( + (julianday(end_time) - julianday(start_time)) * 24 * 60, 1 + ) as duration_minutes + FROM batches + ORDER BY end_time DESC + LIMIT ? + """, + (limit,), + ) + + batches = [] + for row in cur.fetchall(): + batches.append( + { + "date": row["counting_date"], + "batch_number": row["batch_number"], + "count": row["count"], + "start_time": row["start_time"], + "end_time": row["end_time"], + "duration_minutes": row["duration_minutes"] or 0, + } + ) + + conn.close() + return jsonify(batches) + + +@app.route("/api/available-dates") +def api_available_dates(): + """Get list of all available counting dates.""" + conn = get_db() + cur = conn.cursor() + + cur.execute(""" + SELECT counting_date, total_count, total_batches + FROM daily_summaries + ORDER BY counting_date DESC + """) + + dates = [] + for row in cur.fetchall(): + dates.append( + { + "date": row["counting_date"], + "total_count": row["total_count"], + "total_batches": row["total_batches"], + } + ) + + conn.close() + return jsonify(dates) + + +# ------------------------------------------------------------------ # +# CSV Export Routes +# ------------------------------------------------------------------ # +@app.route("/api/export-daily-csv") +def export_daily_csv(): + """Export daily summary records as CSV.""" + days = request.args.get("days", 30, type=int) + date_from = (datetime.now() - timedelta(days=days)).date().isoformat() + + conn = get_db() + cur = conn.cursor() + + cur.execute( + """ + SELECT counting_date, total_count, total_batches, + ROUND(CAST(total_count AS FLOAT) / NULLIF(total_batches, 0), 1) as avg_per_batch + FROM daily_summaries + WHERE counting_date >= ? + ORDER BY counting_date ASC + """, + (date_from,), + ) + + output = StringIO() + writer = csv.writer(output) + writer.writerow(["Date", "Total Count", "Total Batches", "Avg per Batch"]) + + for row in cur.fetchall(): + writer.writerow( + [ + row["counting_date"], + row["total_count"], + row["total_batches"], + row["avg_per_batch"] or 0, + ] + ) + + conn.close() + + csv_data = output.getvalue() + output.close() + + filename = f"daily_records_{datetime.now().strftime('%Y%m%d_%H%M%S')}.csv" + return Response( + csv_data, + mimetype="text/csv", + headers={ + "Content-Disposition": f"attachment; filename={filename}", + "Content-Type": "text/csv; charset=utf-8", + }, + ) + + +@app.route("/api/export-day-csv/") +def export_day_csv(date): + """Export batch details for a specific day as CSV.""" + conn = get_db() + cur = conn.cursor() + + # Batches for this day + cur.execute( + """ + SELECT batch_number, count, start_time, end_time, + ROUND( + (julianday(end_time) - julianday(start_time)) * 24 * 60, 1 + ) as duration_minutes + FROM batches + WHERE counting_date = ? + ORDER BY batch_number ASC + """, + (date,), + ) + + output = StringIO() + writer = csv.writer(output) + writer.writerow( + ["Batch Number", "Count", "Start Time", "End Time", "Duration (min)"] + ) + + for row in cur.fetchall(): + writer.writerow( + [ + row["batch_number"], + row["count"], + row["start_time"], + row["end_time"], + row["duration_minutes"] or 0, + ] + ) + + conn.close() + + csv_data = output.getvalue() + output.close() + + filename = f"day_detail_{date}.csv" + return Response( + csv_data, + mimetype="text/csv", + headers={ + "Content-Disposition": f"attachment; filename={filename}", + "Content-Type": "text/csv; charset=utf-8", + }, + ) + + +# ------------------------------------------------------------------ # +# Run +# ------------------------------------------------------------------ # +if __name__ == "__main__": + # Suppress Flask's default request log spam + WSGIRequestHandler.protocol_version = "HTTP/1.1" + + port = int(os.getenv("DASHBOARD_PORT", 80)) + host = os.getenv("DASHBOARD_HOST", "0.0.0.0") + debug = os.getenv("FLASK_DEBUG", "false").lower() == "true" + + print(f"🚀 Dashboard running at http://{host}:{port}") + app.run(host=host, port=port, debug=debug) diff --git a/counter_service.py b/counter_service.py deleted file mode 100644 index efaa55f..0000000 --- a/counter_service.py +++ /dev/null @@ -1,580 +0,0 @@ -#!/usr/bin/env python3 -""" -Frigate Object Batch Counter Service - -Listens to Frigate MQTT events, counts unique objects per batch, -and persists results to SQLite when a batch ends. -""" - -import os -import sys -import json -import sqlite3 -import threading -import time -import logging -import signal -from datetime import datetime, timedelta -from pathlib import Path - -import paho.mqtt.client as mqtt - - -class FrigateCounterService: - def __init__(self): - self.setup_logging() - self.load_config() - self.init_db() - self.state_lock = threading.Lock() - self.batch_timer = None - self.shutdown_event = threading.Event() - self.current_state = self.load_state() - - # Previous state - self.previous_state = self.current_state - - # Telenan Batch Label - self.ignore_batch_label = False - self.ignore_batch_label_timer = None - self.sleep_after_batch_label_detected = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) - # Sleep none blocking - self.sleep_after_batch_label_timeout = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) - self.sleep_after_batch_label_timer = None - self.sleep_after_batch_label = False - # ------------------------------------------------------------------ # - # Setup & Config - # ------------------------------------------------------------------ # - def setup_logging(self): - level = getattr(logging, os.getenv("LOG_LEVEL", "INFO").upper(), logging.INFO) - logging.basicConfig( - level=level, - format="%(asctime)s [%(levelname)s] %(message)s", - handlers=[logging.StreamHandler(sys.stdout)], - ) - self.logger = logging.getLogger(__name__) - - def load_config(self): - self.mqtt_host = os.getenv("FRIGATE_MQTT_HOST", "localhost") - self.mqtt_port = int(os.getenv("FRIGATE_MQTT_PORT", "1883")) - self.mqtt_user = os.getenv("FRIGATE_MQTT_USER") - self.mqtt_pass = os.getenv("FRIGATE_MQTT_PASS") - self.mqtt_topic = os.getenv("FRIGATE_MQTT_TOPIC", "frigate/events") - - self.camera_name = os.getenv("CAMERA_NAME") - if not self.camera_name: - raise ValueError("Environment variable CAMERA_NAME is required") - - self.object_label = os.getenv("OBJECT_LABEL", "ayam-potong") - self.batch_label = os.getenv("BATCH_LABEL", "telenan") # 20260514 - Adding Label telenan for new batch sign - self.ignore_batch_label_timeout = float(os.getenv("IGNORE_BATCH_LABEL_TIMEOUT_SECONDS", "30")) - self.batch_timeout = float(os.getenv("BATCH_TIMEOUT_SECONDS", "300")) - self.cutoff_time_str = os.getenv("DAILY_CUTOFF_TIME", "17:00") - - # Validate cutoff format HH:MM - datetime.strptime(self.cutoff_time_str, "%H:%M") - - self.db_path = os.getenv("DB_PATH", "frigate_counter.db") - self.state_file = os.getenv("STATE_FILE", "current_batch.json") - - self.min_duration_per_batch = int(os.getenv("MIN_DURATION_PER_BATCH", "60")) - self.min_object_per_batch = int(os.getenv("MIN_OBJECT_PER_BATCH", "60")) - # ------------------------------------------------------------------ # - # Database - # ------------------------------------------------------------------ # - def init_db(self): - self.db = sqlite3.connect(self.db_path, check_same_thread=False) - cur = self.db.cursor() - - cur.execute( - """ - CREATE TABLE IF NOT EXISTS batches ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - counting_date TEXT NOT NULL, - batch_number INTEGER NOT NULL, - camera_name TEXT NOT NULL, - object_label TEXT NOT NULL, - count INTEGER NOT NULL, - start_time TEXT NOT NULL, - end_time TEXT NOT NULL, - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, - UNIQUE(counting_date, batch_number, camera_name, object_label) - ) - """ - ) - - cur.execute( - """ - CREATE TABLE IF NOT EXISTS daily_summaries ( - 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_batches INTEGER NOT NULL DEFAULT 0, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, - UNIQUE(counting_date, camera_name, object_label) - ) - """ - ) - - self.db.commit() - - # ------------------------------------------------------------------ # - # Counting-date logic (day ends at cutoff, e.g. 17:00) - # ------------------------------------------------------------------ # - def get_counting_date(self, dt=None): - """Return the business-day string that ends at cutoff_time.""" - if dt is None: - dt = datetime.now() - cutoff = datetime.strptime(self.cutoff_time_str, "%H:%M").time() - # e.g. cutoff 17:00 => 16:59 belongs to today, 17:00 belongs to tomorrow - if dt.time() < cutoff: - #if dt.time() >= cutoff: - return dt.date().isoformat() - return (dt.date() + timedelta(days=1)).isoformat() - #return (dt.date() - timedelta(days=1)).isoformat() - - # ------------------------------------------------------------------ # - # State persistence (JSON) – survives restarts - # ------------------------------------------------------------------ # - def load_state(self): - if not Path(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.logger.warning( - "State file belongs to previous counting day (%s). " - "Finalizing it before starting fresh.", - state.get("counting_date"), - ) - self._insert_batch( - state["counting_date"], - state["batch_number"], - state["count"], - state["start_time"], - datetime.now().isoformat(), - ) - Path(self.state_file).unlink(missing_ok=True) - return None - - self.logger.info( - "Resumed batch #%s from %s with count=%s", - state["batch_number"], - state["start_time"], - state["count"], - ) - # Restart the inactivity timer - """ - Comment this to remove reset timer - """ - self._reset_batch_timer() - return state - - except Exception as exc: - self.logger.error("Failed to load state file: %s", exc) - return None - - def save_state(self): - if self.current_state is None: - Path(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) - - # ------------------------------------------------------------------ # - # Batch lifecycle - # ------------------------------------------------------------------ # - def get_next_batch_number(self, counting_date): - cur = self.db.cursor() - cur.execute( - """ - SELECT COALESCE(MAX(batch_number), 0) - FROM batches - WHERE counting_date = ? AND camera_name = ? AND object_label = ? - """, - (counting_date, self.camera_name, self.object_label), - ) - return cur.fetchone()[0] + 1 - - def start_new_batch(self, counting_date): - batch_number = self.get_next_batch_number(counting_date) - now = datetime.now().isoformat() - - """ START - Carry last 50 previous counted_ids into next batch """ - counted_ids = [] - if self.previous_state is not None: - try: - counted_ids = self.previous_state["counted_event_ids"][-50:] - except: - counted_ids = [] - - # Print - #self.logger.info(counted_ids) - - self.current_state = { - "counting_date": counting_date, - "batch_number": batch_number, - "count": 0, - "start_time": now, - "last_detection_time": now, - "counted_event_ids": counted_ids, - } - """ END - Carry last 50 previous counted_ids into next batch """ - - """ - self.current_state = { - "counting_date": counting_date, - "batch_number": batch_number, - "count": 0, - "start_time": now, - "last_detection_time": now, - "counted_event_ids": [], - } - """ - self.save_state() - self.logger.info( - "Started batch #%s for %s (%s)", - batch_number, - counting_date, - self.object_label, - ) - - def process_detection(self, event_id, label): - """ - Called for every matching Frigate event (new/update/end). - Deduplicates by event_id and resets the 5-minute batch timer. - """ - should_reset_timer = False - - if label == self.object_label: - with self.state_lock: - counting_date = self.get_counting_date() - - # 1. No active batch -> start one - if self.current_state is None: - self.start_new_batch(counting_date) - should_reset_timer = True - - # 2. Cutoff crossed since batch started -> finalize old, start new - elif self.current_state["counting_date"] != counting_date: - self._end_batch_locked() - self.start_new_batch(counting_date) - should_reset_timer = True - - # 3. Deduplicate event ID - if event_id not in self.current_state["counted_event_ids"]: - self.current_state["count"] += 1 - self.current_state["counted_event_ids"].append(event_id) - self.logger.info( - "Counted %s (event %s) | batch #%s total: %s", - self.object_label, - event_id, - self.current_state["batch_number"], - self.current_state["count"], - ) - - # Always refresh last_detection_time so the batch stays alive - self.current_state["last_detection_time"] = datetime.now().isoformat() - self.save_state() - should_reset_timer = True - - elif label == self.batch_label: - self._ignore_batch_label() - self.end_batch() - - # Sleep Blocking - #self.logger.info("Batch Label detected. Sleep for %s seconds.", self.sleep_after_batch_label_detected) - #time.sleep(self.sleep_after_batch_label_detected) - - # Sleep Non Blocking - #self._sleep_after_batch_label() - - should_reset_timer = True - - """ - Comment this to remove reset timer - """ - if should_reset_timer: - self._reset_batch_timer() - - def _sleep_after_batch_label(self): - if not self.sleep_after_batch_label_timer: - self.sleep_after_batch_label = True - self.sleep_after_batch_label_timer = threading.Timer(self.sleep_after_batch_label_timeout, self._on_sleep_after_batch_label_timeout) - self.sleep_after_batch_label_timer.daemon = True - self.sleep_after_batch_label_timer.start() - self.logger.info("Sleep (non blocking) after Batch Label for %ss. self.sleep_after_batch_label = %s", self.sleep_after_batch_label_timeout, self.sleep_after_batch_label) - - def _on_sleep_after_batch_label_timeout(self): - self.sleep_after_batch_label_timer = None - self.sleep_after_batch_label = False - self.logger.info("Sleep (non blocking) after Batch Label is done. self.sleep_after_batch_label = %s", self.sleep_after_batch_label) - - def _ignore_batch_label(self): - if not self.ignore_batch_label_timer: - self.ignore_batch_label = True - self.ignore_batch_label_timer = threading.Timer(self.ignore_batch_label_timeout, self._on_ignore_batch_label_timeout) - self.ignore_batch_label_timer.daemon = True - self.ignore_batch_label_timer.start() - self.logger.info("Ignore Batch Label for %ss. self.ignore_batch_label = %s", self.ignore_batch_label_timeout, self.ignore_batch_label) - - def _reset_batch_timer(self): - if self.batch_timer: - self.batch_timer.cancel() - self.batch_timer = threading.Timer(self.batch_timeout, self._on_batch_timeout) - self.batch_timer.daemon = True - self.batch_timer.start() - - def _on_ignore_batch_label_timeout(self): - self.ignore_batch_label_timer = None - self.ignore_batch_label = False - self.logger.info("Ignore Batch Label is done. self.ignore_batch_label = %s", self.ignore_batch_label) - #self.end_batch() - - def _on_batch_timeout(self): - self.logger.info("Batch inactivity timeout (%ss) reached", self.batch_timeout) - self.end_batch() - - def end_batch(self): - with self.state_lock: - self._end_batch_locked() - - def _end_batch_locked(self): - if self.current_state is None: - return - - self.previous_state = self.current_state - state = self.current_state - - """ Check Duration """ - start_time_obj = datetime.fromisoformat(state["start_time"]) - end_time_obj = datetime.now() - duration = end_time_obj - start_time_obj - duration_seconds = duration.total_seconds() - - # Minimal conditions per batch - #if state["count"] == 0 or duration_seconds < self.min_duration_per_batch: - if state["count"] < self.min_object_per_batch or duration_seconds < self.min_duration_per_batch: - # Nothing to persist - self.current_state = None - self.save_state() - - """ - Related to _reset_batch_timer - """ - if self.batch_timer: - self.batch_timer.cancel() - self.batch_timer = None - - return - - end_time = datetime.now().isoformat() - - count_per_second = state["count"] / duration_seconds - - try: - self._insert_batch( - state["counting_date"], - state["batch_number"], - state["count"], - state["start_time"], - end_time, - ) - self.logger.info( - "Batch #%s ended | count=%s | duration=%s | cps=%s", - state["batch_number"], - state["count"], - duration, - count_per_second - ) - except Exception as exc: - self.logger.error("Failed to persist batch: %s", exc) - # Leave state intact so we can retry on next timeout - return - - self.current_state = None - self.save_state() - """ - Related to _reset_batch_timer - """ - if self.batch_timer: - self.batch_timer.cancel() - self.batch_timer = None - - def _insert_batch(self, counting_date, batch_number, count, start_time, end_time): - """Atomic insert into batches + upsert daily summary.""" - cur = self.db.cursor() - - cur.execute( - """ - INSERT INTO batches - (counting_date, batch_number, camera_name, object_label, count, start_time, end_time) - VALUES (?, ?, ?, ?, ?, ?, ?) - """, - ( - counting_date, - batch_number, - self.camera_name, - self.object_label, - count, - start_time, - end_time, - ), - ) - - cur.execute( - """ - INSERT INTO daily_summaries - (counting_date, camera_name, object_label, total_count, total_batches) - VALUES (?, ?, ?, ?, 1) - ON CONFLICT(counting_date, camera_name, object_label) - DO UPDATE SET - total_count = total_count + excluded.total_count, - total_batches = total_batches + excluded.total_batches, - updated_at = CURRENT_TIMESTAMP - """, - (counting_date, self.camera_name, self.object_label, count), - ) - - self.db.commit() - - # Log running totals for the day - cur.execute( - """ - SELECT total_count, total_batches - FROM daily_summaries - WHERE counting_date = ? AND camera_name = ? AND object_label = ? - """, - (counting_date, self.camera_name, self.object_label), - ) - row = cur.fetchone() - if row: - self.logger.info( - "Daily totals for %s: %s objects across %s batch(es)", - counting_date, - row[0], - row[1], - ) - - # ------------------------------------------------------------------ # - # Cutoff watcher (forces batch end at 17:00 etc.) - # ------------------------------------------------------------------ # - def cutoff_watcher(self): - """Runs every minute to force-close a batch when the business day rolls over.""" - 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.logger.info("Daily cutoff reached – finalizing batch") - self._end_batch_locked() - - # ------------------------------------------------------------------ # - # MQTT callbacks - # ------------------------------------------------------------------ # - def on_connect(self, client, userdata, flags, rc): - if rc == 0: - self.logger.info("MQTT connected to %s:%s", self.mqtt_host, self.mqtt_port) - client.subscribe(self.mqtt_topic) - self.logger.info("Subscribed to %s", self.mqtt_topic) - else: - self.logger.error("MQTT connection failed, code=%s", rc) - - def on_message(self, client, userdata, msg): - - # Sleep Non Blocking after batch label detected - if self.sleep_after_batch_label: - return - - try: - payload = json.loads(msg.payload.decode("utf-8")) - after = payload.get("after", {}) - - - """ Start - 20260530 for ignore stationary object """ - stationary = after.get("stationary") - if stationary: - return - """ End - 20260530 for ignore stationary object """ - - """ 20260514 For Telenan """ - label = after.get("label") - camera = after.get("camera") - - if camera != self.camera_name: - return - - #if after.get("label") != self.object_label: - if label != self.object_label: - if label == self.batch_label and self.ignore_batch_label: - return - - #if label != self.batch_label and self.ignore_batch_label: - # return - - event_id = after.get("id") - if not event_id: - return - - self.process_detection(event_id, label) - - except json.JSONDecodeError: - self.logger.warning("Received non-JSON payload on %s", msg.topic) - except Exception as exc: - self.logger.exception("Error handling MQTT message: %s", exc) - - # ------------------------------------------------------------------ # - # Run / Shutdown - # ------------------------------------------------------------------ # - def run(self): - # Graceful shutdown on SIGINT / SIGTERM - signal.signal(signal.SIGINT, lambda _s, _f: self.shutdown()) - signal.signal(signal.SIGTERM, lambda _s, _f: self.shutdown()) - - # Support both paho-mqtt v1 and v2 - try: - self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION1) - except (AttributeError, TypeError): - self.client = mqtt.Client() - - if self.mqtt_user and self.mqtt_pass: - self.client.username_pw_set(self.mqtt_user, self.mqtt_pass) - - self.client.on_connect = self.on_connect - self.client.on_message = self.on_message - - # Start background cutoff watcher - watcher = threading.Thread(target=self.cutoff_watcher, daemon=True) - watcher.start() - - try: - self.client.connect(self.mqtt_host, self.mqtt_port, keepalive=60) - self.client.loop_forever() - except Exception as exc: - self.logger.error("MQTT loop error: %s", exc) - finally: - self.shutdown() - - def shutdown(self): - if self.shutdown_event.is_set(): - return - self.logger.info("Shutting down...") - self.shutdown_event.set() - try: - self.client.disconnect() - except Exception: - pass - self.end_batch() - self.db.close() - self.logger.info("Shutdown complete") - - -if __name__ == "__main__": - service = FrigateCounterService() - service.run() diff --git a/counter_service.py b/counter_service.py new file mode 120000 index 0000000..b18d9a7 --- /dev/null +++ b/counter_service.py @@ -0,0 +1 @@ +counter_service.py-zone_check \ No newline at end of file diff --git a/counter_service.py-zone_check b/counter_service.py-zone_check new file mode 100644 index 0000000..64beab4 --- /dev/null +++ b/counter_service.py-zone_check @@ -0,0 +1,585 @@ +#!/usr/bin/env python3 +""" +Frigate Object Batch Counter Service + +Listens to Frigate MQTT events, counts unique objects per batch, +and persists results to SQLite when a batch ends. +""" + +import os +import sys +import json +import sqlite3 +import threading +import time +import logging +import signal +from datetime import datetime, timedelta +from pathlib import Path + +import paho.mqtt.client as mqtt + + +class FrigateCounterService: + def __init__(self): + self.setup_logging() + self.load_config() + self.init_db() + self.state_lock = threading.Lock() + self.batch_timer = None + self.shutdown_event = threading.Event() + self.current_state = self.load_state() + + # Previous state + self.previous_state = self.current_state + + # Telenan Batch Label + self.ignore_batch_label = False + self.ignore_batch_label_timer = None + self.sleep_after_batch_label_detected = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + # Sleep none blocking + self.sleep_after_batch_label_timeout = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + + # Counter Zone + self.zone_counter = os.getenv("ZONE_COUNTER", "zone_counter") + # ------------------------------------------------------------------ # + # Setup & Config + # ------------------------------------------------------------------ # + def setup_logging(self): + level = getattr(logging, os.getenv("LOG_LEVEL", "INFO").upper(), logging.INFO) + logging.basicConfig( + level=level, + format="%(asctime)s [%(levelname)s] %(message)s", + handlers=[logging.StreamHandler(sys.stdout)], + ) + self.logger = logging.getLogger(__name__) + + def load_config(self): + self.mqtt_host = os.getenv("FRIGATE_MQTT_HOST", "localhost") + self.mqtt_port = int(os.getenv("FRIGATE_MQTT_PORT", "1883")) + self.mqtt_user = os.getenv("FRIGATE_MQTT_USER") + self.mqtt_pass = os.getenv("FRIGATE_MQTT_PASS") + self.mqtt_topic = os.getenv("FRIGATE_MQTT_TOPIC", "frigate/events") + + self.camera_name = os.getenv("CAMERA_NAME") + if not self.camera_name: + raise ValueError("Environment variable CAMERA_NAME is required") + + self.object_label = os.getenv("OBJECT_LABEL", "ayam-potong") + self.batch_label = os.getenv("BATCH_LABEL", "telenan") # 20260514 - Adding Label telenan for new batch sign + self.ignore_batch_label_timeout = float(os.getenv("IGNORE_BATCH_LABEL_TIMEOUT_SECONDS", "30")) + self.batch_timeout = float(os.getenv("BATCH_TIMEOUT_SECONDS", "300")) + self.cutoff_time_str = os.getenv("DAILY_CUTOFF_TIME", "17:00") + + # Validate cutoff format HH:MM + datetime.strptime(self.cutoff_time_str, "%H:%M") + + self.db_path = os.getenv("DB_PATH", "frigate_counter.db") + self.state_file = os.getenv("STATE_FILE", "current_batch.json") + + self.min_duration_per_batch = int(os.getenv("MIN_DURATION_PER_BATCH", "60")) + self.min_object_per_batch = int(os.getenv("MIN_OBJECT_PER_BATCH", "60")) + # ------------------------------------------------------------------ # + # Database + # ------------------------------------------------------------------ # + def init_db(self): + self.db = sqlite3.connect(self.db_path, check_same_thread=False) + cur = self.db.cursor() + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS batches ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + counting_date TEXT NOT NULL, + batch_number INTEGER NOT NULL, + camera_name TEXT NOT NULL, + object_label TEXT NOT NULL, + count INTEGER NOT NULL, + start_time TEXT NOT NULL, + end_time TEXT NOT NULL, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, batch_number, camera_name, object_label) + ) + """ + ) + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS daily_summaries ( + 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_batches INTEGER NOT NULL DEFAULT 0, + updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, camera_name, object_label) + ) + """ + ) + + self.db.commit() + + # ------------------------------------------------------------------ # + # Counting-date logic (day ends at cutoff, e.g. 17:00) + # ------------------------------------------------------------------ # + def get_counting_date(self, dt=None): + """Return the business-day string that ends at cutoff_time.""" + if dt is None: + dt = datetime.now() + cutoff = datetime.strptime(self.cutoff_time_str, "%H:%M").time() + # e.g. cutoff 17:00 => 16:59 belongs to today, 17:00 belongs to tomorrow + if dt.time() < cutoff: + #if dt.time() >= cutoff: + return dt.date().isoformat() + return (dt.date() + timedelta(days=1)).isoformat() + #return (dt.date() - timedelta(days=1)).isoformat() + + # ------------------------------------------------------------------ # + # State persistence (JSON) – survives restarts + # ------------------------------------------------------------------ # + def load_state(self): + if not Path(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.logger.warning( + "State file belongs to previous counting day (%s). " + "Finalizing it before starting fresh.", + state.get("counting_date"), + ) + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + datetime.now().isoformat(), + ) + Path(self.state_file).unlink(missing_ok=True) + return None + + self.logger.info( + "Resumed batch #%s from %s with count=%s", + state["batch_number"], + state["start_time"], + state["count"], + ) + # Restart the inactivity timer + """ + Comment this to remove reset timer + """ + self._reset_batch_timer() + return state + + except Exception as exc: + self.logger.error("Failed to load state file: %s", exc) + return None + + def save_state(self): + if self.current_state is None: + Path(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) + + # ------------------------------------------------------------------ # + # Batch lifecycle + # ------------------------------------------------------------------ # + def get_next_batch_number(self, counting_date): + cur = self.db.cursor() + cur.execute( + """ + SELECT COALESCE(MAX(batch_number), 0) + FROM batches + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + return cur.fetchone()[0] + 1 + + def start_new_batch(self, counting_date): + batch_number = self.get_next_batch_number(counting_date) + now = datetime.now().isoformat() + + """ START - Carry last 50 previous counted_ids into next batch """ + counted_ids = [] + if self.previous_state is not None: + try: + counted_ids = self.previous_state["counted_event_ids"][-50:] + except: + counted_ids = [] + + # Print + #self.logger.info(counted_ids) + + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": counted_ids, + } + """ END - Carry last 50 previous counted_ids into next batch """ + + """ + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": [], + } + """ + self.save_state() + self.logger.info( + "Started batch #%s for %s (%s)", + batch_number, + counting_date, + self.object_label, + ) + + def process_detection(self, event_id, label): + """ + Called for every matching Frigate event (new/update/end). + Deduplicates by event_id and resets the 5-minute batch timer. + """ + should_reset_timer = False + + if label == self.object_label: + with self.state_lock: + counting_date = self.get_counting_date() + + # 1. No active batch -> start one + if self.current_state is None: + self.start_new_batch(counting_date) + should_reset_timer = True + + # 2. Cutoff crossed since batch started -> finalize old, start new + elif self.current_state["counting_date"] != counting_date: + self._end_batch_locked() + self.start_new_batch(counting_date) + should_reset_timer = True + + # 3. Deduplicate event ID + if event_id not in self.current_state["counted_event_ids"]: + self.current_state["count"] += 1 + self.current_state["counted_event_ids"].append(event_id) + self.logger.info( + "Counted %s (event %s) in zone %s | batch #%s total: %s", + self.object_label, + event_id, + self.zone_counter, + self.current_state["batch_number"], + self.current_state["count"], + ) + + # Always refresh last_detection_time so the batch stays alive + self.current_state["last_detection_time"] = datetime.now().isoformat() + self.save_state() + should_reset_timer = True + + elif label == self.batch_label: + self._ignore_batch_label() + self.end_batch() + + # Sleep Blocking + #self.logger.info("Batch Label detected. Sleep for %s seconds.", self.sleep_after_batch_label_detected) + #time.sleep(self.sleep_after_batch_label_detected) + + # Sleep Non Blocking + #self._sleep_after_batch_label() + + should_reset_timer = True + + """ + Comment this to remove reset timer + """ + if should_reset_timer: + self._reset_batch_timer() + + def _sleep_after_batch_label(self): + if not self.sleep_after_batch_label_timer: + self.sleep_after_batch_label = True + self.sleep_after_batch_label_timer = threading.Timer(self.sleep_after_batch_label_timeout, self._on_sleep_after_batch_label_timeout) + self.sleep_after_batch_label_timer.daemon = True + self.sleep_after_batch_label_timer.start() + self.logger.info("Sleep (non blocking) after Batch Label for %ss. self.sleep_after_batch_label = %s", self.sleep_after_batch_label_timeout, self.sleep_after_batch_label) + + def _on_sleep_after_batch_label_timeout(self): + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + self.logger.info("Sleep (non blocking) after Batch Label is done. self.sleep_after_batch_label = %s", self.sleep_after_batch_label) + + def _ignore_batch_label(self): + if not self.ignore_batch_label_timer: + self.ignore_batch_label = True + self.ignore_batch_label_timer = threading.Timer(self.ignore_batch_label_timeout, self._on_ignore_batch_label_timeout) + self.ignore_batch_label_timer.daemon = True + self.ignore_batch_label_timer.start() + self.logger.info("Ignore Batch Label for %ss. self.ignore_batch_label = %s", self.ignore_batch_label_timeout, self.ignore_batch_label) + + def _reset_batch_timer(self): + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = threading.Timer(self.batch_timeout, self._on_batch_timeout) + self.batch_timer.daemon = True + self.batch_timer.start() + + def _on_ignore_batch_label_timeout(self): + self.ignore_batch_label_timer = None + self.ignore_batch_label = False + self.logger.info("Ignore Batch Label is done. self.ignore_batch_label = %s", self.ignore_batch_label) + #self.end_batch() + + def _on_batch_timeout(self): + self.logger.info("Batch inactivity timeout (%ss) reached", self.batch_timeout) + self.end_batch() + + def end_batch(self): + with self.state_lock: + self._end_batch_locked() + + def _end_batch_locked(self): + if self.current_state is None: + return + + self.previous_state = self.current_state + state = self.current_state + + """ Check Duration """ + start_time_obj = datetime.fromisoformat(state["start_time"]) + end_time_obj = datetime.now() + duration = end_time_obj - start_time_obj + duration_seconds = duration.total_seconds() + + # Minimal conditions per batch + #if state["count"] == 0 or duration_seconds < self.min_duration_per_batch: + if state["count"] < self.min_object_per_batch or duration_seconds < self.min_duration_per_batch: + # Nothing to persist + self.current_state = None + self.save_state() + + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + return + + end_time = datetime.now().isoformat() + + count_per_second = state["count"] / duration_seconds + + try: + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + end_time, + ) + self.logger.info( + "Batch #%s ended | count=%s | duration=%s | cps=%s", + state["batch_number"], + state["count"], + duration, + count_per_second + ) + except Exception as exc: + self.logger.error("Failed to persist batch: %s", exc) + # Leave state intact so we can retry on next timeout + return + + self.current_state = None + self.save_state() + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + def _insert_batch(self, counting_date, batch_number, count, start_time, end_time): + """Atomic insert into batches + upsert daily summary.""" + cur = self.db.cursor() + + cur.execute( + """ + INSERT INTO batches + (counting_date, batch_number, camera_name, object_label, count, start_time, end_time) + VALUES (?, ?, ?, ?, ?, ?, ?) + """, + ( + counting_date, + batch_number, + self.camera_name, + self.object_label, + count, + start_time, + end_time, + ), + ) + + cur.execute( + """ + INSERT INTO daily_summaries + (counting_date, camera_name, object_label, total_count, total_batches) + VALUES (?, ?, ?, ?, 1) + ON CONFLICT(counting_date, camera_name, object_label) + DO UPDATE SET + total_count = total_count + excluded.total_count, + total_batches = total_batches + excluded.total_batches, + updated_at = CURRENT_TIMESTAMP + """, + (counting_date, self.camera_name, self.object_label, count), + ) + + self.db.commit() + + # Log running totals for the day + cur.execute( + """ + SELECT total_count, total_batches + FROM daily_summaries + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + row = cur.fetchone() + if row: + self.logger.info( + "Daily totals for %s: %s objects across %s batch(es)", + counting_date, + row[0], + row[1], + ) + + # ------------------------------------------------------------------ # + # Cutoff watcher (forces batch end at 17:00 etc.) + # ------------------------------------------------------------------ # + def cutoff_watcher(self): + """Runs every minute to force-close a batch when the business day rolls over.""" + 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.logger.info("Daily cutoff reached – finalizing batch") + self._end_batch_locked() + + # ------------------------------------------------------------------ # + # MQTT callbacks + # ------------------------------------------------------------------ # + def on_connect(self, client, userdata, flags, rc): + if rc == 0: + self.logger.info("MQTT connected to %s:%s", self.mqtt_host, self.mqtt_port) + client.subscribe(self.mqtt_topic) + self.logger.info("Subscribed to %s", self.mqtt_topic) + else: + self.logger.error("MQTT connection failed, code=%s", rc) + + def on_message(self, client, userdata, msg): + + #self.logger.info(msg) + # Sleep Non Blocking after batch label detected + #if self.sleep_after_batch_label: + # return + + try: + payload = json.loads(msg.payload.decode("utf-8")) + after = payload.get("after", {}) + + + """ Start - 20260530 for ignore stationary object """ + stationary = after.get("stationary") + if stationary: + return + """ End - 20260530 for ignore stationary object """ + + """ 20260514 For Telenan """ + camera = after.get("camera") + if camera != self.camera_name: + return + + label = after.get("label") + #if after.get("label") != self.object_label: + if label != self.object_label: + if label == self.batch_label and self.ignore_batch_label: + return + + event_id = after.get("id") + if not event_id: + return + + zones = after.get("entered_zones") + if self.zone_counter not in zones: + return + + self.process_detection(event_id, label) + + except json.JSONDecodeError: + self.logger.warning("Received non-JSON payload on %s", msg.topic) + except Exception as exc: + self.logger.exception("Error handling MQTT message: %s", exc) + + # ------------------------------------------------------------------ # + # Run / Shutdown + # ------------------------------------------------------------------ # + def run(self): + # Graceful shutdown on SIGINT / SIGTERM + signal.signal(signal.SIGINT, lambda _s, _f: self.shutdown()) + signal.signal(signal.SIGTERM, lambda _s, _f: self.shutdown()) + + # Support both paho-mqtt v1 and v2 + try: + self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION1) + except (AttributeError, TypeError): + self.client = mqtt.Client() + + if self.mqtt_user and self.mqtt_pass: + self.client.username_pw_set(self.mqtt_user, self.mqtt_pass) + + self.client.on_connect = self.on_connect + self.client.on_message = self.on_message + + # Start background cutoff watcher + watcher = threading.Thread(target=self.cutoff_watcher, daemon=True) + watcher.start() + + try: + self.client.connect(self.mqtt_host, self.mqtt_port, keepalive=60) + self.client.loop_forever() + except Exception as exc: + self.logger.error("MQTT loop error: %s", exc) + finally: + self.shutdown() + + def shutdown(self): + if self.shutdown_event.is_set(): + return + self.logger.info("Shutting down...") + self.shutdown_event.set() + try: + self.client.disconnect() + except Exception: + pass + self.end_batch() + self.db.close() + self.logger.info("Shutdown complete") + + +if __name__ == "__main__": + service = FrigateCounterService() + service.run() diff --git a/counter_service.py.20260604 b/counter_service.py.20260604 new file mode 100644 index 0000000..8790e5e --- /dev/null +++ b/counter_service.py.20260604 @@ -0,0 +1,526 @@ +#!/usr/bin/env python3 +""" +Frigate Object Batch Counter Service + +Listens to Frigate MQTT events, counts unique objects per batch, +and persists results to SQLite when a batch ends. +""" + +import os +import sys +import json +import sqlite3 +import threading +import time +import logging +import signal +from datetime import datetime, timedelta +from pathlib import Path + +import paho.mqtt.client as mqtt + + +class FrigateCounterService: + def __init__(self): + self.setup_logging() + self.load_config() + self.init_db() + self.state_lock = threading.Lock() + self.batch_timer = None + self.shutdown_event = threading.Event() + self.current_state = self.load_state() + + # Telenan Batch Label + self.ignore_batch_label = False + self.ignore_batch_label_timer = None + self.sleep_after_batch_label_detected = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + # ------------------------------------------------------------------ # + # Setup & Config + # ------------------------------------------------------------------ # + def setup_logging(self): + level = getattr(logging, os.getenv("LOG_LEVEL", "INFO").upper(), logging.INFO) + logging.basicConfig( + level=level, + format="%(asctime)s [%(levelname)s] %(message)s", + handlers=[logging.StreamHandler(sys.stdout)], + ) + self.logger = logging.getLogger(__name__) + + def load_config(self): + self.mqtt_host = os.getenv("FRIGATE_MQTT_HOST", "localhost") + self.mqtt_port = int(os.getenv("FRIGATE_MQTT_PORT", "1883")) + self.mqtt_user = os.getenv("FRIGATE_MQTT_USER") + self.mqtt_pass = os.getenv("FRIGATE_MQTT_PASS") + self.mqtt_topic = os.getenv("FRIGATE_MQTT_TOPIC", "frigate/events") + + self.camera_name = os.getenv("CAMERA_NAME") + if not self.camera_name: + raise ValueError("Environment variable CAMERA_NAME is required") + + self.object_label = os.getenv("OBJECT_LABEL", "ayam-potong") + self.batch_label = os.getenv("BATCH_LABEL", "telenan") # 20260514 - Adding Label telenan for new batch sign + self.ignore_batch_label_timeout = float(os.getenv("IGNORE_BATCH_LABEL_TIMEOUT_SECONDS", "30")) + self.batch_timeout = float(os.getenv("BATCH_TIMEOUT_SECONDS", "300")) + self.cutoff_time_str = os.getenv("DAILY_CUTOFF_TIME", "17:00") + + # Validate cutoff format HH:MM + datetime.strptime(self.cutoff_time_str, "%H:%M") + + self.db_path = os.getenv("DB_PATH", "frigate_counter.db") + self.state_file = os.getenv("STATE_FILE", "current_batch.json") + + self.min_duration_per_batch = int(os.getenv("MIN_DURATION_PER_BATCH", "60")) + self.min_object_per_batch = int(os.getenv("MIN_OBJECT_PER_BATCH", "60")) + # ------------------------------------------------------------------ # + # Database + # ------------------------------------------------------------------ # + def init_db(self): + self.db = sqlite3.connect(self.db_path, check_same_thread=False) + cur = self.db.cursor() + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS batches ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + counting_date TEXT NOT NULL, + batch_number INTEGER NOT NULL, + camera_name TEXT NOT NULL, + object_label TEXT NOT NULL, + count INTEGER NOT NULL, + start_time TEXT NOT NULL, + end_time TEXT NOT NULL, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, batch_number, camera_name, object_label) + ) + """ + ) + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS daily_summaries ( + 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_batches INTEGER NOT NULL DEFAULT 0, + updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, camera_name, object_label) + ) + """ + ) + + self.db.commit() + + # ------------------------------------------------------------------ # + # Counting-date logic (day ends at cutoff, e.g. 17:00) + # ------------------------------------------------------------------ # + def get_counting_date(self, dt=None): + """Return the business-day string that ends at cutoff_time.""" + if dt is None: + dt = datetime.now() + cutoff = datetime.strptime(self.cutoff_time_str, "%H:%M").time() + # e.g. cutoff 17:00 => 16:59 belongs to today, 17:00 belongs to tomorrow + if dt.time() < cutoff: + #if dt.time() >= cutoff: + return dt.date().isoformat() + return (dt.date() + timedelta(days=1)).isoformat() + #return (dt.date() - timedelta(days=1)).isoformat() + + # ------------------------------------------------------------------ # + # State persistence (JSON) – survives restarts + # ------------------------------------------------------------------ # + def load_state(self): + if not Path(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.logger.warning( + "State file belongs to previous counting day (%s). " + "Finalizing it before starting fresh.", + state.get("counting_date"), + ) + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + datetime.now().isoformat(), + ) + Path(self.state_file).unlink(missing_ok=True) + return None + + self.logger.info( + "Resumed batch #%s from %s with count=%s", + state["batch_number"], + state["start_time"], + state["count"], + ) + # Restart the inactivity timer + """ + Comment this to remove reset timer + """ + self._reset_batch_timer() + return state + + except Exception as exc: + self.logger.error("Failed to load state file: %s", exc) + return None + + def save_state(self): + if self.current_state is None: + Path(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) + + # ------------------------------------------------------------------ # + # Batch lifecycle + # ------------------------------------------------------------------ # + def get_next_batch_number(self, counting_date): + cur = self.db.cursor() + cur.execute( + """ + SELECT COALESCE(MAX(batch_number), 0) + FROM batches + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + return cur.fetchone()[0] + 1 + + def start_new_batch(self, counting_date): + batch_number = self.get_next_batch_number(counting_date) + now = datetime.now().isoformat() + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": [], + } + self.save_state() + self.logger.info( + "Started batch #%s for %s (%s)", + batch_number, + counting_date, + self.object_label, + ) + + def process_detection(self, event_id, label): + """ + Called for every matching Frigate event (new/update/end). + Deduplicates by event_id and resets the 5-minute batch timer. + """ + should_reset_timer = False + + if label == self.object_label: + with self.state_lock: + counting_date = self.get_counting_date() + + # 1. No active batch -> start one + if self.current_state is None: + self.start_new_batch(counting_date) + should_reset_timer = True + + # 2. Cutoff crossed since batch started -> finalize old, start new + elif self.current_state["counting_date"] != counting_date: + self._end_batch_locked() + self.start_new_batch(counting_date) + should_reset_timer = True + + # 3. Deduplicate event ID + if event_id not in self.current_state["counted_event_ids"]: + self.current_state["count"] += 1 + self.current_state["counted_event_ids"].append(event_id) + self.logger.info( + "Counted %s (event %s) | batch #%s total: %s", + self.object_label, + event_id, + self.current_state["batch_number"], + self.current_state["count"], + ) + + # Always refresh last_detection_time so the batch stays alive + self.current_state["last_detection_time"] = datetime.now().isoformat() + self.save_state() + should_reset_timer = True + + elif label == self.batch_label: + self._ignore_batch_label() + self.end_batch() + + self.logger.info("Batch Label detected. Sleep for %s seconds.", self.sleep_after_batch_label_detected) + time.sleep(self.sleep_after_batch_label_detected) + + should_reset_timer = True + + """ + Comment this to remove reset timer + """ + if should_reset_timer: + self._reset_batch_timer() + + def _ignore_batch_label(self): + if not self.ignore_batch_label_timer: + self.ignore_batch_label = True + self.ignore_batch_label_timer = threading.Timer(self.ignore_batch_label_timeout, self._on_ignore_batch_label_timeout) + self.ignore_batch_label_timer.daemon = True + self.ignore_batch_label_timer.start() + self.logger.info("Ignore Batch Label for %ss. self.ignore_batch_label = %s", self.ignore_batch_label_timeout, self.ignore_batch_label) + + def _reset_batch_timer(self): + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = threading.Timer(self.batch_timeout, self._on_batch_timeout) + self.batch_timer.daemon = True + self.batch_timer.start() + + def _on_ignore_batch_label_timeout(self): + self.ignore_batch_label_timer = None + self.ignore_batch_label = False + self.logger.info("Ignore Batch Label is done. self.ignore_batch_label = %s", self.ignore_batch_label) + #self.end_batch() + + def _on_batch_timeout(self): + self.logger.info("Batch inactivity timeout (%ss) reached", self.batch_timeout) + self.end_batch() + + def end_batch(self): + with self.state_lock: + self._end_batch_locked() + + def _end_batch_locked(self): + if self.current_state is None: + return + + state = self.current_state + + """ Check Duration """ + start_time_obj = datetime.fromisoformat(state["start_time"]) + end_time_obj = datetime.now() + duration = end_time_obj - start_time_obj + duration_seconds = duration.total_seconds() + + # Minimal conditions per batch + #if state["count"] == 0 or duration_seconds < self.min_duration_per_batch: + if state["count"] < self.min_object_per_batch or duration_seconds < self.min_duration_per_batch: + # Nothing to persist + self.current_state = None + self.save_state() + + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + return + + end_time = datetime.now().isoformat() + + count_per_second = state["count"] / duration_seconds + + try: + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + end_time, + ) + self.logger.info( + "Batch #%s ended | count=%s | duration=%s | cps=%s", + state["batch_number"], + state["count"], + duration, + count_per_second + ) + except Exception as exc: + self.logger.error("Failed to persist batch: %s", exc) + # Leave state intact so we can retry on next timeout + return + + self.current_state = None + self.save_state() + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + def _insert_batch(self, counting_date, batch_number, count, start_time, end_time): + """Atomic insert into batches + upsert daily summary.""" + cur = self.db.cursor() + + cur.execute( + """ + INSERT INTO batches + (counting_date, batch_number, camera_name, object_label, count, start_time, end_time) + VALUES (?, ?, ?, ?, ?, ?, ?) + """, + ( + counting_date, + batch_number, + self.camera_name, + self.object_label, + count, + start_time, + end_time, + ), + ) + + cur.execute( + """ + INSERT INTO daily_summaries + (counting_date, camera_name, object_label, total_count, total_batches) + VALUES (?, ?, ?, ?, 1) + ON CONFLICT(counting_date, camera_name, object_label) + DO UPDATE SET + total_count = total_count + excluded.total_count, + total_batches = total_batches + excluded.total_batches, + updated_at = CURRENT_TIMESTAMP + """, + (counting_date, self.camera_name, self.object_label, count), + ) + + self.db.commit() + + # Log running totals for the day + cur.execute( + """ + SELECT total_count, total_batches + FROM daily_summaries + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + row = cur.fetchone() + if row: + self.logger.info( + "Daily totals for %s: %s objects across %s batch(es)", + counting_date, + row[0], + row[1], + ) + + # ------------------------------------------------------------------ # + # Cutoff watcher (forces batch end at 17:00 etc.) + # ------------------------------------------------------------------ # + def cutoff_watcher(self): + """Runs every minute to force-close a batch when the business day rolls over.""" + 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.logger.info("Daily cutoff reached – finalizing batch") + self._end_batch_locked() + + # ------------------------------------------------------------------ # + # MQTT callbacks + # ------------------------------------------------------------------ # + def on_connect(self, client, userdata, flags, rc): + if rc == 0: + self.logger.info("MQTT connected to %s:%s", self.mqtt_host, self.mqtt_port) + client.subscribe(self.mqtt_topic) + self.logger.info("Subscribed to %s", self.mqtt_topic) + else: + self.logger.error("MQTT connection failed, code=%s", rc) + + def on_message(self, client, userdata, msg): + try: + payload = json.loads(msg.payload.decode("utf-8")) + after = payload.get("after", {}) + + + """ Start - 20260530 for ignore stationary object """ + stationary = after.get("stationary") + if stationary: + return + """ End - 20260530 for ignore stationary object """ + + """ 20260514 For Telenan """ + label = after.get("label") + camera = after.get("camera") + + if camera != self.camera_name: + return + + #if after.get("label") != self.object_label: + if label != self.object_label: + if label == self.batch_label and self.ignore_batch_label: + return + + #if label != self.batch_label and self.ignore_batch_label: + # return + + event_id = after.get("id") + if not event_id: + return + + self.process_detection(event_id, label) + + except json.JSONDecodeError: + self.logger.warning("Received non-JSON payload on %s", msg.topic) + except Exception as exc: + self.logger.exception("Error handling MQTT message: %s", exc) + + # ------------------------------------------------------------------ # + # Run / Shutdown + # ------------------------------------------------------------------ # + def run(self): + # Graceful shutdown on SIGINT / SIGTERM + signal.signal(signal.SIGINT, lambda _s, _f: self.shutdown()) + signal.signal(signal.SIGTERM, lambda _s, _f: self.shutdown()) + + # Support both paho-mqtt v1 and v2 + try: + self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION1) + except (AttributeError, TypeError): + self.client = mqtt.Client() + + if self.mqtt_user and self.mqtt_pass: + self.client.username_pw_set(self.mqtt_user, self.mqtt_pass) + + self.client.on_connect = self.on_connect + self.client.on_message = self.on_message + + # Start background cutoff watcher + watcher = threading.Thread(target=self.cutoff_watcher, daemon=True) + watcher.start() + + try: + self.client.connect(self.mqtt_host, self.mqtt_port, keepalive=60) + self.client.loop_forever() + except Exception as exc: + self.logger.error("MQTT loop error: %s", exc) + finally: + self.shutdown() + + def shutdown(self): + if self.shutdown_event.is_set(): + return + self.logger.info("Shutting down...") + self.shutdown_event.set() + try: + self.client.disconnect() + except Exception: + pass + self.end_batch() + self.db.close() + self.logger.info("Shutdown complete") + + +if __name__ == "__main__": + service = FrigateCounterService() + service.run() diff --git a/counter_service.py.20260609 b/counter_service.py.20260609 new file mode 100644 index 0000000..93fcc70 --- /dev/null +++ b/counter_service.py.20260609 @@ -0,0 +1,552 @@ +#!/usr/bin/env python3 +""" +Frigate Object Batch Counter Service + +Listens to Frigate MQTT events, counts unique objects per batch, +and persists results to SQLite when a batch ends. +""" + +import os +import sys +import json +import sqlite3 +import threading +import time +import logging +import signal +from datetime import datetime, timedelta +from pathlib import Path + +import paho.mqtt.client as mqtt + + +class FrigateCounterService: + def __init__(self): + self.setup_logging() + self.load_config() + self.init_db() + self.state_lock = threading.Lock() + self.batch_timer = None + self.shutdown_event = threading.Event() + self.current_state = self.load_state() + + # Telenan Batch Label + self.ignore_batch_label = False + self.ignore_batch_label_timer = None + self.sleep_after_batch_label_detected = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + # Sleep none blocking + self.sleep_after_batch_label_timeout = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + # ------------------------------------------------------------------ # + # Setup & Config + # ------------------------------------------------------------------ # + def setup_logging(self): + level = getattr(logging, os.getenv("LOG_LEVEL", "INFO").upper(), logging.INFO) + logging.basicConfig( + level=level, + format="%(asctime)s [%(levelname)s] %(message)s", + handlers=[logging.StreamHandler(sys.stdout)], + ) + self.logger = logging.getLogger(__name__) + + def load_config(self): + self.mqtt_host = os.getenv("FRIGATE_MQTT_HOST", "localhost") + self.mqtt_port = int(os.getenv("FRIGATE_MQTT_PORT", "1883")) + self.mqtt_user = os.getenv("FRIGATE_MQTT_USER") + self.mqtt_pass = os.getenv("FRIGATE_MQTT_PASS") + self.mqtt_topic = os.getenv("FRIGATE_MQTT_TOPIC", "frigate/events") + + self.camera_name = os.getenv("CAMERA_NAME") + if not self.camera_name: + raise ValueError("Environment variable CAMERA_NAME is required") + + self.object_label = os.getenv("OBJECT_LABEL", "ayam-potong") + self.batch_label = os.getenv("BATCH_LABEL", "telenan") # 20260514 - Adding Label telenan for new batch sign + self.ignore_batch_label_timeout = float(os.getenv("IGNORE_BATCH_LABEL_TIMEOUT_SECONDS", "30")) + self.batch_timeout = float(os.getenv("BATCH_TIMEOUT_SECONDS", "300")) + self.cutoff_time_str = os.getenv("DAILY_CUTOFF_TIME", "17:00") + + # Validate cutoff format HH:MM + datetime.strptime(self.cutoff_time_str, "%H:%M") + + self.db_path = os.getenv("DB_PATH", "frigate_counter.db") + self.state_file = os.getenv("STATE_FILE", "current_batch.json") + + self.min_duration_per_batch = int(os.getenv("MIN_DURATION_PER_BATCH", "60")) + self.min_object_per_batch = int(os.getenv("MIN_OBJECT_PER_BATCH", "60")) + # ------------------------------------------------------------------ # + # Database + # ------------------------------------------------------------------ # + def init_db(self): + self.db = sqlite3.connect(self.db_path, check_same_thread=False) + cur = self.db.cursor() + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS batches ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + counting_date TEXT NOT NULL, + batch_number INTEGER NOT NULL, + camera_name TEXT NOT NULL, + object_label TEXT NOT NULL, + count INTEGER NOT NULL, + start_time TEXT NOT NULL, + end_time TEXT NOT NULL, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, batch_number, camera_name, object_label) + ) + """ + ) + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS daily_summaries ( + 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_batches INTEGER NOT NULL DEFAULT 0, + updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, camera_name, object_label) + ) + """ + ) + + self.db.commit() + + # ------------------------------------------------------------------ # + # Counting-date logic (day ends at cutoff, e.g. 17:00) + # ------------------------------------------------------------------ # + def get_counting_date(self, dt=None): + """Return the business-day string that ends at cutoff_time.""" + if dt is None: + dt = datetime.now() + cutoff = datetime.strptime(self.cutoff_time_str, "%H:%M").time() + # e.g. cutoff 17:00 => 16:59 belongs to today, 17:00 belongs to tomorrow + if dt.time() < cutoff: + #if dt.time() >= cutoff: + return dt.date().isoformat() + return (dt.date() + timedelta(days=1)).isoformat() + #return (dt.date() - timedelta(days=1)).isoformat() + + # ------------------------------------------------------------------ # + # State persistence (JSON) – survives restarts + # ------------------------------------------------------------------ # + def load_state(self): + if not Path(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.logger.warning( + "State file belongs to previous counting day (%s). " + "Finalizing it before starting fresh.", + state.get("counting_date"), + ) + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + datetime.now().isoformat(), + ) + Path(self.state_file).unlink(missing_ok=True) + return None + + self.logger.info( + "Resumed batch #%s from %s with count=%s", + state["batch_number"], + state["start_time"], + state["count"], + ) + # Restart the inactivity timer + """ + Comment this to remove reset timer + """ + self._reset_batch_timer() + return state + + except Exception as exc: + self.logger.error("Failed to load state file: %s", exc) + return None + + def save_state(self): + if self.current_state is None: + Path(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) + + # ------------------------------------------------------------------ # + # Batch lifecycle + # ------------------------------------------------------------------ # + def get_next_batch_number(self, counting_date): + cur = self.db.cursor() + cur.execute( + """ + SELECT COALESCE(MAX(batch_number), 0) + FROM batches + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + return cur.fetchone()[0] + 1 + + def start_new_batch(self, counting_date): + batch_number = self.get_next_batch_number(counting_date) + now = datetime.now().isoformat() + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": [], + } + self.save_state() + self.logger.info( + "Started batch #%s for %s (%s)", + batch_number, + counting_date, + self.object_label, + ) + + def process_detection(self, event_id, label): + """ + Called for every matching Frigate event (new/update/end). + Deduplicates by event_id and resets the 5-minute batch timer. + """ + should_reset_timer = False + + if label == self.object_label: + with self.state_lock: + counting_date = self.get_counting_date() + + # 1. No active batch -> start one + if self.current_state is None: + self.start_new_batch(counting_date) + should_reset_timer = True + + # 2. Cutoff crossed since batch started -> finalize old, start new + elif self.current_state["counting_date"] != counting_date: + self._end_batch_locked() + self.start_new_batch(counting_date) + should_reset_timer = True + + # 3. Deduplicate event ID + if event_id not in self.current_state["counted_event_ids"]: + self.current_state["count"] += 1 + self.current_state["counted_event_ids"].append(event_id) + self.logger.info( + "Counted %s (event %s) | batch #%s total: %s", + self.object_label, + event_id, + self.current_state["batch_number"], + self.current_state["count"], + ) + + # Always refresh last_detection_time so the batch stays alive + self.current_state["last_detection_time"] = datetime.now().isoformat() + self.save_state() + should_reset_timer = True + + elif label == self.batch_label: + self._ignore_batch_label() + self.end_batch() + + # Sleep Blocking + #self.logger.info("Batch Label detected. Sleep for %s seconds.", self.sleep_after_batch_label_detected) + #time.sleep(self.sleep_after_batch_label_detected) + + # Sleep Non Blocking + self._sleep_after_batch_label() + + should_reset_timer = True + + """ + Comment this to remove reset timer + """ + if should_reset_timer: + self._reset_batch_timer() + + def _sleep_after_batch_label(self): + if not self.sleep_after_batch_label_timer: + self.sleep_after_batch_label = True + self.sleep_after_batch_label_timer = threading.Timer(self.sleep_after_batch_label_timeout, self._on_sleep_after_batch_label_timeout) + self.sleep_after_batch_label_timer.daemon = True + self.sleep_after_batch_label_timer.start() + self.logger.info("Sleep (non blocking) after Batch Label for %ss. self.sleep_after_batch_label = %s", self.sleep_after_batch_label_timeout, self.sleep_after_batch_label) + + def _on_sleep_after_batch_label_timeout(self): + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + self.logger.info("Sleep (non blocking) after Batch Label is done. self.sleep_after_batch_label = %s", self.sleep_after_batch_label) + + def _ignore_batch_label(self): + if not self.ignore_batch_label_timer: + self.ignore_batch_label = True + self.ignore_batch_label_timer = threading.Timer(self.ignore_batch_label_timeout, self._on_ignore_batch_label_timeout) + self.ignore_batch_label_timer.daemon = True + self.ignore_batch_label_timer.start() + self.logger.info("Ignore Batch Label for %ss. self.ignore_batch_label = %s", self.ignore_batch_label_timeout, self.ignore_batch_label) + + def _reset_batch_timer(self): + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = threading.Timer(self.batch_timeout, self._on_batch_timeout) + self.batch_timer.daemon = True + self.batch_timer.start() + + def _on_ignore_batch_label_timeout(self): + self.ignore_batch_label_timer = None + self.ignore_batch_label = False + self.logger.info("Ignore Batch Label is done. self.ignore_batch_label = %s", self.ignore_batch_label) + #self.end_batch() + + def _on_batch_timeout(self): + self.logger.info("Batch inactivity timeout (%ss) reached", self.batch_timeout) + self.end_batch() + + def end_batch(self): + with self.state_lock: + self._end_batch_locked() + + def _end_batch_locked(self): + if self.current_state is None: + return + + state = self.current_state + + """ Check Duration """ + start_time_obj = datetime.fromisoformat(state["start_time"]) + end_time_obj = datetime.now() + duration = end_time_obj - start_time_obj + duration_seconds = duration.total_seconds() + + # Minimal conditions per batch + #if state["count"] == 0 or duration_seconds < self.min_duration_per_batch: + if state["count"] < self.min_object_per_batch or duration_seconds < self.min_duration_per_batch: + # Nothing to persist + self.current_state = None + self.save_state() + + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + return + + end_time = datetime.now().isoformat() + + count_per_second = state["count"] / duration_seconds + + try: + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + end_time, + ) + self.logger.info( + "Batch #%s ended | count=%s | duration=%s | cps=%s", + state["batch_number"], + state["count"], + duration, + count_per_second + ) + except Exception as exc: + self.logger.error("Failed to persist batch: %s", exc) + # Leave state intact so we can retry on next timeout + return + + self.current_state = None + self.save_state() + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + def _insert_batch(self, counting_date, batch_number, count, start_time, end_time): + """Atomic insert into batches + upsert daily summary.""" + cur = self.db.cursor() + + cur.execute( + """ + INSERT INTO batches + (counting_date, batch_number, camera_name, object_label, count, start_time, end_time) + VALUES (?, ?, ?, ?, ?, ?, ?) + """, + ( + counting_date, + batch_number, + self.camera_name, + self.object_label, + count, + start_time, + end_time, + ), + ) + + cur.execute( + """ + INSERT INTO daily_summaries + (counting_date, camera_name, object_label, total_count, total_batches) + VALUES (?, ?, ?, ?, 1) + ON CONFLICT(counting_date, camera_name, object_label) + DO UPDATE SET + total_count = total_count + excluded.total_count, + total_batches = total_batches + excluded.total_batches, + updated_at = CURRENT_TIMESTAMP + """, + (counting_date, self.camera_name, self.object_label, count), + ) + + self.db.commit() + + # Log running totals for the day + cur.execute( + """ + SELECT total_count, total_batches + FROM daily_summaries + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + row = cur.fetchone() + if row: + self.logger.info( + "Daily totals for %s: %s objects across %s batch(es)", + counting_date, + row[0], + row[1], + ) + + # ------------------------------------------------------------------ # + # Cutoff watcher (forces batch end at 17:00 etc.) + # ------------------------------------------------------------------ # + def cutoff_watcher(self): + """Runs every minute to force-close a batch when the business day rolls over.""" + 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.logger.info("Daily cutoff reached – finalizing batch") + self._end_batch_locked() + + # ------------------------------------------------------------------ # + # MQTT callbacks + # ------------------------------------------------------------------ # + def on_connect(self, client, userdata, flags, rc): + if rc == 0: + self.logger.info("MQTT connected to %s:%s", self.mqtt_host, self.mqtt_port) + client.subscribe(self.mqtt_topic) + self.logger.info("Subscribed to %s", self.mqtt_topic) + else: + self.logger.error("MQTT connection failed, code=%s", rc) + + def on_message(self, client, userdata, msg): + + # Sleep Non Blocking after batch label detected + if self.sleep_after_batch_label: + return + + try: + payload = json.loads(msg.payload.decode("utf-8")) + after = payload.get("after", {}) + + + """ Start - 20260530 for ignore stationary object """ + stationary = after.get("stationary") + if stationary: + return + """ End - 20260530 for ignore stationary object """ + + """ 20260514 For Telenan """ + label = after.get("label") + camera = after.get("camera") + + if camera != self.camera_name: + return + + #if after.get("label") != self.object_label: + if label != self.object_label: + if label == self.batch_label and self.ignore_batch_label: + return + + #if label != self.batch_label and self.ignore_batch_label: + # return + + event_id = after.get("id") + if not event_id: + return + + self.process_detection(event_id, label) + + except json.JSONDecodeError: + self.logger.warning("Received non-JSON payload on %s", msg.topic) + except Exception as exc: + self.logger.exception("Error handling MQTT message: %s", exc) + + # ------------------------------------------------------------------ # + # Run / Shutdown + # ------------------------------------------------------------------ # + def run(self): + # Graceful shutdown on SIGINT / SIGTERM + signal.signal(signal.SIGINT, lambda _s, _f: self.shutdown()) + signal.signal(signal.SIGTERM, lambda _s, _f: self.shutdown()) + + # Support both paho-mqtt v1 and v2 + try: + self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION1) + except (AttributeError, TypeError): + self.client = mqtt.Client() + + if self.mqtt_user and self.mqtt_pass: + self.client.username_pw_set(self.mqtt_user, self.mqtt_pass) + + self.client.on_connect = self.on_connect + self.client.on_message = self.on_message + + # Start background cutoff watcher + watcher = threading.Thread(target=self.cutoff_watcher, daemon=True) + watcher.start() + + try: + self.client.connect(self.mqtt_host, self.mqtt_port, keepalive=60) + self.client.loop_forever() + except Exception as exc: + self.logger.error("MQTT loop error: %s", exc) + finally: + self.shutdown() + + def shutdown(self): + if self.shutdown_event.is_set(): + return + self.logger.info("Shutting down...") + self.shutdown_event.set() + try: + self.client.disconnect() + except Exception: + pass + self.end_batch() + self.db.close() + self.logger.info("Shutdown complete") + + +if __name__ == "__main__": + service = FrigateCounterService() + service.run() diff --git a/counter_service.py.20260612 b/counter_service.py.20260612 new file mode 100644 index 0000000..efaa55f --- /dev/null +++ b/counter_service.py.20260612 @@ -0,0 +1,580 @@ +#!/usr/bin/env python3 +""" +Frigate Object Batch Counter Service + +Listens to Frigate MQTT events, counts unique objects per batch, +and persists results to SQLite when a batch ends. +""" + +import os +import sys +import json +import sqlite3 +import threading +import time +import logging +import signal +from datetime import datetime, timedelta +from pathlib import Path + +import paho.mqtt.client as mqtt + + +class FrigateCounterService: + def __init__(self): + self.setup_logging() + self.load_config() + self.init_db() + self.state_lock = threading.Lock() + self.batch_timer = None + self.shutdown_event = threading.Event() + self.current_state = self.load_state() + + # Previous state + self.previous_state = self.current_state + + # Telenan Batch Label + self.ignore_batch_label = False + self.ignore_batch_label_timer = None + self.sleep_after_batch_label_detected = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + # Sleep none blocking + self.sleep_after_batch_label_timeout = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + # ------------------------------------------------------------------ # + # Setup & Config + # ------------------------------------------------------------------ # + def setup_logging(self): + level = getattr(logging, os.getenv("LOG_LEVEL", "INFO").upper(), logging.INFO) + logging.basicConfig( + level=level, + format="%(asctime)s [%(levelname)s] %(message)s", + handlers=[logging.StreamHandler(sys.stdout)], + ) + self.logger = logging.getLogger(__name__) + + def load_config(self): + self.mqtt_host = os.getenv("FRIGATE_MQTT_HOST", "localhost") + self.mqtt_port = int(os.getenv("FRIGATE_MQTT_PORT", "1883")) + self.mqtt_user = os.getenv("FRIGATE_MQTT_USER") + self.mqtt_pass = os.getenv("FRIGATE_MQTT_PASS") + self.mqtt_topic = os.getenv("FRIGATE_MQTT_TOPIC", "frigate/events") + + self.camera_name = os.getenv("CAMERA_NAME") + if not self.camera_name: + raise ValueError("Environment variable CAMERA_NAME is required") + + self.object_label = os.getenv("OBJECT_LABEL", "ayam-potong") + self.batch_label = os.getenv("BATCH_LABEL", "telenan") # 20260514 - Adding Label telenan for new batch sign + self.ignore_batch_label_timeout = float(os.getenv("IGNORE_BATCH_LABEL_TIMEOUT_SECONDS", "30")) + self.batch_timeout = float(os.getenv("BATCH_TIMEOUT_SECONDS", "300")) + self.cutoff_time_str = os.getenv("DAILY_CUTOFF_TIME", "17:00") + + # Validate cutoff format HH:MM + datetime.strptime(self.cutoff_time_str, "%H:%M") + + self.db_path = os.getenv("DB_PATH", "frigate_counter.db") + self.state_file = os.getenv("STATE_FILE", "current_batch.json") + + self.min_duration_per_batch = int(os.getenv("MIN_DURATION_PER_BATCH", "60")) + self.min_object_per_batch = int(os.getenv("MIN_OBJECT_PER_BATCH", "60")) + # ------------------------------------------------------------------ # + # Database + # ------------------------------------------------------------------ # + def init_db(self): + self.db = sqlite3.connect(self.db_path, check_same_thread=False) + cur = self.db.cursor() + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS batches ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + counting_date TEXT NOT NULL, + batch_number INTEGER NOT NULL, + camera_name TEXT NOT NULL, + object_label TEXT NOT NULL, + count INTEGER NOT NULL, + start_time TEXT NOT NULL, + end_time TEXT NOT NULL, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, batch_number, camera_name, object_label) + ) + """ + ) + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS daily_summaries ( + 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_batches INTEGER NOT NULL DEFAULT 0, + updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, camera_name, object_label) + ) + """ + ) + + self.db.commit() + + # ------------------------------------------------------------------ # + # Counting-date logic (day ends at cutoff, e.g. 17:00) + # ------------------------------------------------------------------ # + def get_counting_date(self, dt=None): + """Return the business-day string that ends at cutoff_time.""" + if dt is None: + dt = datetime.now() + cutoff = datetime.strptime(self.cutoff_time_str, "%H:%M").time() + # e.g. cutoff 17:00 => 16:59 belongs to today, 17:00 belongs to tomorrow + if dt.time() < cutoff: + #if dt.time() >= cutoff: + return dt.date().isoformat() + return (dt.date() + timedelta(days=1)).isoformat() + #return (dt.date() - timedelta(days=1)).isoformat() + + # ------------------------------------------------------------------ # + # State persistence (JSON) – survives restarts + # ------------------------------------------------------------------ # + def load_state(self): + if not Path(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.logger.warning( + "State file belongs to previous counting day (%s). " + "Finalizing it before starting fresh.", + state.get("counting_date"), + ) + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + datetime.now().isoformat(), + ) + Path(self.state_file).unlink(missing_ok=True) + return None + + self.logger.info( + "Resumed batch #%s from %s with count=%s", + state["batch_number"], + state["start_time"], + state["count"], + ) + # Restart the inactivity timer + """ + Comment this to remove reset timer + """ + self._reset_batch_timer() + return state + + except Exception as exc: + self.logger.error("Failed to load state file: %s", exc) + return None + + def save_state(self): + if self.current_state is None: + Path(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) + + # ------------------------------------------------------------------ # + # Batch lifecycle + # ------------------------------------------------------------------ # + def get_next_batch_number(self, counting_date): + cur = self.db.cursor() + cur.execute( + """ + SELECT COALESCE(MAX(batch_number), 0) + FROM batches + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + return cur.fetchone()[0] + 1 + + def start_new_batch(self, counting_date): + batch_number = self.get_next_batch_number(counting_date) + now = datetime.now().isoformat() + + """ START - Carry last 50 previous counted_ids into next batch """ + counted_ids = [] + if self.previous_state is not None: + try: + counted_ids = self.previous_state["counted_event_ids"][-50:] + except: + counted_ids = [] + + # Print + #self.logger.info(counted_ids) + + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": counted_ids, + } + """ END - Carry last 50 previous counted_ids into next batch """ + + """ + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": [], + } + """ + self.save_state() + self.logger.info( + "Started batch #%s for %s (%s)", + batch_number, + counting_date, + self.object_label, + ) + + def process_detection(self, event_id, label): + """ + Called for every matching Frigate event (new/update/end). + Deduplicates by event_id and resets the 5-minute batch timer. + """ + should_reset_timer = False + + if label == self.object_label: + with self.state_lock: + counting_date = self.get_counting_date() + + # 1. No active batch -> start one + if self.current_state is None: + self.start_new_batch(counting_date) + should_reset_timer = True + + # 2. Cutoff crossed since batch started -> finalize old, start new + elif self.current_state["counting_date"] != counting_date: + self._end_batch_locked() + self.start_new_batch(counting_date) + should_reset_timer = True + + # 3. Deduplicate event ID + if event_id not in self.current_state["counted_event_ids"]: + self.current_state["count"] += 1 + self.current_state["counted_event_ids"].append(event_id) + self.logger.info( + "Counted %s (event %s) | batch #%s total: %s", + self.object_label, + event_id, + self.current_state["batch_number"], + self.current_state["count"], + ) + + # Always refresh last_detection_time so the batch stays alive + self.current_state["last_detection_time"] = datetime.now().isoformat() + self.save_state() + should_reset_timer = True + + elif label == self.batch_label: + self._ignore_batch_label() + self.end_batch() + + # Sleep Blocking + #self.logger.info("Batch Label detected. Sleep for %s seconds.", self.sleep_after_batch_label_detected) + #time.sleep(self.sleep_after_batch_label_detected) + + # Sleep Non Blocking + #self._sleep_after_batch_label() + + should_reset_timer = True + + """ + Comment this to remove reset timer + """ + if should_reset_timer: + self._reset_batch_timer() + + def _sleep_after_batch_label(self): + if not self.sleep_after_batch_label_timer: + self.sleep_after_batch_label = True + self.sleep_after_batch_label_timer = threading.Timer(self.sleep_after_batch_label_timeout, self._on_sleep_after_batch_label_timeout) + self.sleep_after_batch_label_timer.daemon = True + self.sleep_after_batch_label_timer.start() + self.logger.info("Sleep (non blocking) after Batch Label for %ss. self.sleep_after_batch_label = %s", self.sleep_after_batch_label_timeout, self.sleep_after_batch_label) + + def _on_sleep_after_batch_label_timeout(self): + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + self.logger.info("Sleep (non blocking) after Batch Label is done. self.sleep_after_batch_label = %s", self.sleep_after_batch_label) + + def _ignore_batch_label(self): + if not self.ignore_batch_label_timer: + self.ignore_batch_label = True + self.ignore_batch_label_timer = threading.Timer(self.ignore_batch_label_timeout, self._on_ignore_batch_label_timeout) + self.ignore_batch_label_timer.daemon = True + self.ignore_batch_label_timer.start() + self.logger.info("Ignore Batch Label for %ss. self.ignore_batch_label = %s", self.ignore_batch_label_timeout, self.ignore_batch_label) + + def _reset_batch_timer(self): + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = threading.Timer(self.batch_timeout, self._on_batch_timeout) + self.batch_timer.daemon = True + self.batch_timer.start() + + def _on_ignore_batch_label_timeout(self): + self.ignore_batch_label_timer = None + self.ignore_batch_label = False + self.logger.info("Ignore Batch Label is done. self.ignore_batch_label = %s", self.ignore_batch_label) + #self.end_batch() + + def _on_batch_timeout(self): + self.logger.info("Batch inactivity timeout (%ss) reached", self.batch_timeout) + self.end_batch() + + def end_batch(self): + with self.state_lock: + self._end_batch_locked() + + def _end_batch_locked(self): + if self.current_state is None: + return + + self.previous_state = self.current_state + state = self.current_state + + """ Check Duration """ + start_time_obj = datetime.fromisoformat(state["start_time"]) + end_time_obj = datetime.now() + duration = end_time_obj - start_time_obj + duration_seconds = duration.total_seconds() + + # Minimal conditions per batch + #if state["count"] == 0 or duration_seconds < self.min_duration_per_batch: + if state["count"] < self.min_object_per_batch or duration_seconds < self.min_duration_per_batch: + # Nothing to persist + self.current_state = None + self.save_state() + + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + return + + end_time = datetime.now().isoformat() + + count_per_second = state["count"] / duration_seconds + + try: + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + end_time, + ) + self.logger.info( + "Batch #%s ended | count=%s | duration=%s | cps=%s", + state["batch_number"], + state["count"], + duration, + count_per_second + ) + except Exception as exc: + self.logger.error("Failed to persist batch: %s", exc) + # Leave state intact so we can retry on next timeout + return + + self.current_state = None + self.save_state() + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + def _insert_batch(self, counting_date, batch_number, count, start_time, end_time): + """Atomic insert into batches + upsert daily summary.""" + cur = self.db.cursor() + + cur.execute( + """ + INSERT INTO batches + (counting_date, batch_number, camera_name, object_label, count, start_time, end_time) + VALUES (?, ?, ?, ?, ?, ?, ?) + """, + ( + counting_date, + batch_number, + self.camera_name, + self.object_label, + count, + start_time, + end_time, + ), + ) + + cur.execute( + """ + INSERT INTO daily_summaries + (counting_date, camera_name, object_label, total_count, total_batches) + VALUES (?, ?, ?, ?, 1) + ON CONFLICT(counting_date, camera_name, object_label) + DO UPDATE SET + total_count = total_count + excluded.total_count, + total_batches = total_batches + excluded.total_batches, + updated_at = CURRENT_TIMESTAMP + """, + (counting_date, self.camera_name, self.object_label, count), + ) + + self.db.commit() + + # Log running totals for the day + cur.execute( + """ + SELECT total_count, total_batches + FROM daily_summaries + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + row = cur.fetchone() + if row: + self.logger.info( + "Daily totals for %s: %s objects across %s batch(es)", + counting_date, + row[0], + row[1], + ) + + # ------------------------------------------------------------------ # + # Cutoff watcher (forces batch end at 17:00 etc.) + # ------------------------------------------------------------------ # + def cutoff_watcher(self): + """Runs every minute to force-close a batch when the business day rolls over.""" + 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.logger.info("Daily cutoff reached – finalizing batch") + self._end_batch_locked() + + # ------------------------------------------------------------------ # + # MQTT callbacks + # ------------------------------------------------------------------ # + def on_connect(self, client, userdata, flags, rc): + if rc == 0: + self.logger.info("MQTT connected to %s:%s", self.mqtt_host, self.mqtt_port) + client.subscribe(self.mqtt_topic) + self.logger.info("Subscribed to %s", self.mqtt_topic) + else: + self.logger.error("MQTT connection failed, code=%s", rc) + + def on_message(self, client, userdata, msg): + + # Sleep Non Blocking after batch label detected + if self.sleep_after_batch_label: + return + + try: + payload = json.loads(msg.payload.decode("utf-8")) + after = payload.get("after", {}) + + + """ Start - 20260530 for ignore stationary object """ + stationary = after.get("stationary") + if stationary: + return + """ End - 20260530 for ignore stationary object """ + + """ 20260514 For Telenan """ + label = after.get("label") + camera = after.get("camera") + + if camera != self.camera_name: + return + + #if after.get("label") != self.object_label: + if label != self.object_label: + if label == self.batch_label and self.ignore_batch_label: + return + + #if label != self.batch_label and self.ignore_batch_label: + # return + + event_id = after.get("id") + if not event_id: + return + + self.process_detection(event_id, label) + + except json.JSONDecodeError: + self.logger.warning("Received non-JSON payload on %s", msg.topic) + except Exception as exc: + self.logger.exception("Error handling MQTT message: %s", exc) + + # ------------------------------------------------------------------ # + # Run / Shutdown + # ------------------------------------------------------------------ # + def run(self): + # Graceful shutdown on SIGINT / SIGTERM + signal.signal(signal.SIGINT, lambda _s, _f: self.shutdown()) + signal.signal(signal.SIGTERM, lambda _s, _f: self.shutdown()) + + # Support both paho-mqtt v1 and v2 + try: + self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION1) + except (AttributeError, TypeError): + self.client = mqtt.Client() + + if self.mqtt_user and self.mqtt_pass: + self.client.username_pw_set(self.mqtt_user, self.mqtt_pass) + + self.client.on_connect = self.on_connect + self.client.on_message = self.on_message + + # Start background cutoff watcher + watcher = threading.Thread(target=self.cutoff_watcher, daemon=True) + watcher.start() + + try: + self.client.connect(self.mqtt_host, self.mqtt_port, keepalive=60) + self.client.loop_forever() + except Exception as exc: + self.logger.error("MQTT loop error: %s", exc) + finally: + self.shutdown() + + def shutdown(self): + if self.shutdown_event.is_set(): + return + self.logger.info("Shutting down...") + self.shutdown_event.set() + try: + self.client.disconnect() + except Exception: + pass + self.end_batch() + self.db.close() + self.logger.info("Shutdown complete") + + +if __name__ == "__main__": + service = FrigateCounterService() + service.run() diff --git a/counter_service.py.with-line b/counter_service.py.with-line new file mode 100644 index 0000000..df0c85a --- /dev/null +++ b/counter_service.py.with-line @@ -0,0 +1,588 @@ +#!/usr/bin/env python3 +""" +Frigate Object Batch Counter Service + +Listens to Frigate MQTT events, counts unique objects per batch, +and persists results to SQLite when a batch ends. +""" + +import os +import sys +import json +import sqlite3 +import threading +import time +import logging +import signal +from datetime import datetime, timedelta +from pathlib import Path + +import paho.mqtt.client as mqtt + + +class FrigateCounterService: + def __init__(self): + self.setup_logging() + self.load_config() + self.init_db() + self.state_lock = threading.Lock() + self.batch_timer = None + self.shutdown_event = threading.Event() + self.current_state = self.load_state() + + # Previous state + self.previous_state = self.current_state + + # Telenan Batch Label + self.ignore_batch_label = False + self.ignore_batch_label_timer = None + self.sleep_after_batch_label_detected = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + # Sleep none blocking + self.sleep_after_batch_label_timeout = int(os.getenv("SLEEP_AFTER_BATCH_LABEL", 10)) + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + + # Counter Zone + self.zone_counter = os.getenv("ZONE_COUNTER", "zone_counter") + # ------------------------------------------------------------------ # + # Setup & Config + # ------------------------------------------------------------------ # + def setup_logging(self): + level = getattr(logging, os.getenv("LOG_LEVEL", "INFO").upper(), logging.INFO) + logging.basicConfig( + level=level, + format="%(asctime)s [%(levelname)s] %(message)s", + handlers=[logging.StreamHandler(sys.stdout)], + ) + self.logger = logging.getLogger(__name__) + + def load_config(self): + self.mqtt_host = os.getenv("FRIGATE_MQTT_HOST", "localhost") + self.mqtt_port = int(os.getenv("FRIGATE_MQTT_PORT", "1883")) + self.mqtt_user = os.getenv("FRIGATE_MQTT_USER") + self.mqtt_pass = os.getenv("FRIGATE_MQTT_PASS") + self.mqtt_topic = os.getenv("FRIGATE_MQTT_TOPIC", "frigate/events") + + self.camera_name = os.getenv("CAMERA_NAME") + if not self.camera_name: + raise ValueError("Environment variable CAMERA_NAME is required") + + self.object_label = os.getenv("OBJECT_LABEL", "ayam-potong") + self.batch_label = os.getenv("BATCH_LABEL", "telenan") # 20260514 - Adding Label telenan for new batch sign + self.ignore_batch_label_timeout = float(os.getenv("IGNORE_BATCH_LABEL_TIMEOUT_SECONDS", "30")) + self.batch_timeout = float(os.getenv("BATCH_TIMEOUT_SECONDS", "300")) + self.cutoff_time_str = os.getenv("DAILY_CUTOFF_TIME", "17:00") + + # Validate cutoff format HH:MM + datetime.strptime(self.cutoff_time_str, "%H:%M") + + self.db_path = os.getenv("DB_PATH", "frigate_counter.db") + self.state_file = os.getenv("STATE_FILE", "current_batch.json") + + self.min_duration_per_batch = int(os.getenv("MIN_DURATION_PER_BATCH", "60")) + self.min_object_per_batch = int(os.getenv("MIN_OBJECT_PER_BATCH", "60")) + # ------------------------------------------------------------------ # + # Database + # ------------------------------------------------------------------ # + def init_db(self): + self.db = sqlite3.connect(self.db_path, check_same_thread=False) + cur = self.db.cursor() + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS batches ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + counting_date TEXT NOT NULL, + batch_number INTEGER NOT NULL, + camera_name TEXT NOT NULL, + object_label TEXT NOT NULL, + count INTEGER NOT NULL, + start_time TEXT NOT NULL, + end_time TEXT NOT NULL, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, batch_number, camera_name, object_label) + ) + """ + ) + + cur.execute( + """ + CREATE TABLE IF NOT EXISTS daily_summaries ( + 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_batches INTEGER NOT NULL DEFAULT 0, + updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + UNIQUE(counting_date, camera_name, object_label) + ) + """ + ) + + self.db.commit() + + # ------------------------------------------------------------------ # + # Counting-date logic (day ends at cutoff, e.g. 17:00) + # ------------------------------------------------------------------ # + def get_counting_date(self, dt=None): + """Return the business-day string that ends at cutoff_time.""" + if dt is None: + dt = datetime.now() + cutoff = datetime.strptime(self.cutoff_time_str, "%H:%M").time() + # e.g. cutoff 17:00 => 16:59 belongs to today, 17:00 belongs to tomorrow + if dt.time() < cutoff: + #if dt.time() >= cutoff: + return dt.date().isoformat() + return (dt.date() + timedelta(days=1)).isoformat() + #return (dt.date() - timedelta(days=1)).isoformat() + + # ------------------------------------------------------------------ # + # State persistence (JSON) – survives restarts + # ------------------------------------------------------------------ # + def load_state(self): + if not Path(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.logger.warning( + "State file belongs to previous counting day (%s). " + "Finalizing it before starting fresh.", + state.get("counting_date"), + ) + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + datetime.now().isoformat(), + ) + Path(self.state_file).unlink(missing_ok=True) + return None + + self.logger.info( + "Resumed batch #%s from %s with count=%s", + state["batch_number"], + state["start_time"], + state["count"], + ) + # Restart the inactivity timer + """ + Comment this to remove reset timer + """ + self._reset_batch_timer() + return state + + except Exception as exc: + self.logger.error("Failed to load state file: %s", exc) + return None + + def save_state(self): + if self.current_state is None: + Path(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) + + # ------------------------------------------------------------------ # + # Batch lifecycle + # ------------------------------------------------------------------ # + def get_next_batch_number(self, counting_date): + cur = self.db.cursor() + cur.execute( + """ + SELECT COALESCE(MAX(batch_number), 0) + FROM batches + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + return cur.fetchone()[0] + 1 + + def start_new_batch(self, counting_date): + batch_number = self.get_next_batch_number(counting_date) + now = datetime.now().isoformat() + + """ START - Carry last 50 previous counted_ids into next batch """ + counted_ids = [] + if self.previous_state is not None: + try: + counted_ids = self.previous_state["counted_event_ids"][-50:] + except: + counted_ids = [] + + # Print + #self.logger.info(counted_ids) + + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": counted_ids, + } + """ END - Carry last 50 previous counted_ids into next batch """ + + """ + self.current_state = { + "counting_date": counting_date, + "batch_number": batch_number, + "count": 0, + "start_time": now, + "last_detection_time": now, + "counted_event_ids": [], + } + """ + self.save_state() + self.logger.info( + "Started batch #%s for %s (%s)", + batch_number, + counting_date, + self.object_label, + ) + + def process_detection(self, event_id, label): + """ + Called for every matching Frigate event (new/update/end). + Deduplicates by event_id and resets the 5-minute batch timer. + """ + should_reset_timer = False + + if label == self.object_label: + with self.state_lock: + counting_date = self.get_counting_date() + + # 1. No active batch -> start one + if self.current_state is None: + self.start_new_batch(counting_date) + should_reset_timer = True + + # 2. Cutoff crossed since batch started -> finalize old, start new + elif self.current_state["counting_date"] != counting_date: + self._end_batch_locked() + self.start_new_batch(counting_date) + should_reset_timer = True + + # 3. Deduplicate event ID + if event_id not in self.current_state["counted_event_ids"]: + self.current_state["count"] += 1 + self.current_state["counted_event_ids"].append(event_id) + self.logger.info( + "Counted %s (event %s) | batch #%s total: %s", + self.object_label, + event_id, + self.current_state["batch_number"], + self.current_state["count"], + ) + + # Always refresh last_detection_time so the batch stays alive + self.current_state["last_detection_time"] = datetime.now().isoformat() + self.save_state() + should_reset_timer = True + + elif label == self.batch_label: + self._ignore_batch_label() + self.end_batch() + + # Sleep Blocking + #self.logger.info("Batch Label detected. Sleep for %s seconds.", self.sleep_after_batch_label_detected) + #time.sleep(self.sleep_after_batch_label_detected) + + # Sleep Non Blocking + #self._sleep_after_batch_label() + + should_reset_timer = True + + """ + Comment this to remove reset timer + """ + if should_reset_timer: + self._reset_batch_timer() + + def _sleep_after_batch_label(self): + if not self.sleep_after_batch_label_timer: + self.sleep_after_batch_label = True + self.sleep_after_batch_label_timer = threading.Timer(self.sleep_after_batch_label_timeout, self._on_sleep_after_batch_label_timeout) + self.sleep_after_batch_label_timer.daemon = True + self.sleep_after_batch_label_timer.start() + self.logger.info("Sleep (non blocking) after Batch Label for %ss. self.sleep_after_batch_label = %s", self.sleep_after_batch_label_timeout, self.sleep_after_batch_label) + + def _on_sleep_after_batch_label_timeout(self): + self.sleep_after_batch_label_timer = None + self.sleep_after_batch_label = False + self.logger.info("Sleep (non blocking) after Batch Label is done. self.sleep_after_batch_label = %s", self.sleep_after_batch_label) + + def _ignore_batch_label(self): + if not self.ignore_batch_label_timer: + self.ignore_batch_label = True + self.ignore_batch_label_timer = threading.Timer(self.ignore_batch_label_timeout, self._on_ignore_batch_label_timeout) + self.ignore_batch_label_timer.daemon = True + self.ignore_batch_label_timer.start() + self.logger.info("Ignore Batch Label for %ss. self.ignore_batch_label = %s", self.ignore_batch_label_timeout, self.ignore_batch_label) + + def _reset_batch_timer(self): + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = threading.Timer(self.batch_timeout, self._on_batch_timeout) + self.batch_timer.daemon = True + self.batch_timer.start() + + def _on_ignore_batch_label_timeout(self): + self.ignore_batch_label_timer = None + self.ignore_batch_label = False + self.logger.info("Ignore Batch Label is done. self.ignore_batch_label = %s", self.ignore_batch_label) + #self.end_batch() + + def _on_batch_timeout(self): + self.logger.info("Batch inactivity timeout (%ss) reached", self.batch_timeout) + self.end_batch() + + def end_batch(self): + with self.state_lock: + self._end_batch_locked() + + def _end_batch_locked(self): + if self.current_state is None: + return + + self.previous_state = self.current_state + state = self.current_state + + """ Check Duration """ + start_time_obj = datetime.fromisoformat(state["start_time"]) + end_time_obj = datetime.now() + duration = end_time_obj - start_time_obj + duration_seconds = duration.total_seconds() + + # Minimal conditions per batch + #if state["count"] == 0 or duration_seconds < self.min_duration_per_batch: + if state["count"] < self.min_object_per_batch or duration_seconds < self.min_duration_per_batch: + # Nothing to persist + self.current_state = None + self.save_state() + + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + return + + end_time = datetime.now().isoformat() + + count_per_second = state["count"] / duration_seconds + + try: + self._insert_batch( + state["counting_date"], + state["batch_number"], + state["count"], + state["start_time"], + end_time, + ) + self.logger.info( + "Batch #%s ended | count=%s | duration=%s | cps=%s", + state["batch_number"], + state["count"], + duration, + count_per_second + ) + except Exception as exc: + self.logger.error("Failed to persist batch: %s", exc) + # Leave state intact so we can retry on next timeout + return + + self.current_state = None + self.save_state() + """ + Related to _reset_batch_timer + """ + if self.batch_timer: + self.batch_timer.cancel() + self.batch_timer = None + + def _insert_batch(self, counting_date, batch_number, count, start_time, end_time): + """Atomic insert into batches + upsert daily summary.""" + cur = self.db.cursor() + + cur.execute( + """ + INSERT INTO batches + (counting_date, batch_number, camera_name, object_label, count, start_time, end_time) + VALUES (?, ?, ?, ?, ?, ?, ?) + """, + ( + counting_date, + batch_number, + self.camera_name, + self.object_label, + count, + start_time, + end_time, + ), + ) + + cur.execute( + """ + INSERT INTO daily_summaries + (counting_date, camera_name, object_label, total_count, total_batches) + VALUES (?, ?, ?, ?, 1) + ON CONFLICT(counting_date, camera_name, object_label) + DO UPDATE SET + total_count = total_count + excluded.total_count, + total_batches = total_batches + excluded.total_batches, + updated_at = CURRENT_TIMESTAMP + """, + (counting_date, self.camera_name, self.object_label, count), + ) + + self.db.commit() + + # Log running totals for the day + cur.execute( + """ + SELECT total_count, total_batches + FROM daily_summaries + WHERE counting_date = ? AND camera_name = ? AND object_label = ? + """, + (counting_date, self.camera_name, self.object_label), + ) + row = cur.fetchone() + if row: + self.logger.info( + "Daily totals for %s: %s objects across %s batch(es)", + counting_date, + row[0], + row[1], + ) + + # ------------------------------------------------------------------ # + # Cutoff watcher (forces batch end at 17:00 etc.) + # ------------------------------------------------------------------ # + def cutoff_watcher(self): + """Runs every minute to force-close a batch when the business day rolls over.""" + 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.logger.info("Daily cutoff reached – finalizing batch") + self._end_batch_locked() + + # ------------------------------------------------------------------ # + # MQTT callbacks + # ------------------------------------------------------------------ # + def on_connect(self, client, userdata, flags, rc): + if rc == 0: + self.logger.info("MQTT connected to %s:%s", self.mqtt_host, self.mqtt_port) + client.subscribe(self.mqtt_topic) + self.logger.info("Subscribed to %s", self.mqtt_topic) + else: + self.logger.error("MQTT connection failed, code=%s", rc) + + def on_message(self, client, userdata, msg): + + #self.logger.info(msg) + # Sleep Non Blocking after batch label detected + if self.sleep_after_batch_label: + return + + try: + payload = json.loads(msg.payload.decode("utf-8")) + after = payload.get("after", {}) + + + """ Start - 20260530 for ignore stationary object """ + stationary = after.get("stationary") + if stationary: + return + """ End - 20260530 for ignore stationary object """ + + """ 20260514 For Telenan """ + label = after.get("label") + camera = after.get("camera") + + if camera != self.camera_name: + return + + #if after.get("label") != self.object_label: + if label != self.object_label: + if label == self.batch_label and self.ignore_batch_label: + return + + #if label != self.batch_label and self.ignore_batch_label: + # return + + event_id = after.get("id") + if not event_id: + return + + zones = after.get("entered_zones") + if self.zone_counter not in zones: + return + + self.process_detection(event_id, label) + + except json.JSONDecodeError: + self.logger.warning("Received non-JSON payload on %s", msg.topic) + except Exception as exc: + self.logger.exception("Error handling MQTT message: %s", exc) + + # ------------------------------------------------------------------ # + # Run / Shutdown + # ------------------------------------------------------------------ # + def run(self): + # Graceful shutdown on SIGINT / SIGTERM + signal.signal(signal.SIGINT, lambda _s, _f: self.shutdown()) + signal.signal(signal.SIGTERM, lambda _s, _f: self.shutdown()) + + # Support both paho-mqtt v1 and v2 + try: + self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION1) + except (AttributeError, TypeError): + self.client = mqtt.Client() + + if self.mqtt_user and self.mqtt_pass: + self.client.username_pw_set(self.mqtt_user, self.mqtt_pass) + + self.client.on_connect = self.on_connect + self.client.on_message = self.on_message + + # Start background cutoff watcher + watcher = threading.Thread(target=self.cutoff_watcher, daemon=True) + watcher.start() + + try: + self.client.connect(self.mqtt_host, self.mqtt_port, keepalive=60) + self.client.loop_forever() + except Exception as exc: + self.logger.error("MQTT loop error: %s", exc) + finally: + self.shutdown() + + def shutdown(self): + if self.shutdown_event.is_set(): + return + self.logger.info("Shutting down...") + self.shutdown_event.set() + try: + self.client.disconnect() + except Exception: + pass + self.end_batch() + self.db.close() + self.logger.info("Shutdown complete") + + +if __name__ == "__main__": + service = FrigateCounterService() + service.run() diff --git a/frigate_counter.db b/frigate_counter.db index 0f5b5f3..3ee9271 100644 Binary files a/frigate_counter.db and b/frigate_counter.db differ diff --git a/templates/dashboard.html b/templates/dashboard.html index 4bed33d..652a052 100644 --- a/templates/dashboard.html +++ b/templates/dashboard.html @@ -1,796 +1,1045 @@ - - - 🐔 ZenAI APC Dashboard - - - - + * { margin: 0; padding: 0; box-sizing: border-box; } + + :root { + --bg-deep: #06080e; + --bg-surface: rgba(12, 16, 28, 0.9); + --border-glow: rgba(0, 240, 255, 0.12); + --accent: #00f0ff; + --accent2: #7b2fff; + --accent3: #ff2d78; + --accent4: #00ff88; + --text-primary: #dce2f0; + --text-secondary: #6b7394; + --glass: rgba(12, 16, 32, 0.7); + --cell-bg: rgba(255, 255, 255, 0.02); + --cell-border: rgba(255, 255, 255, 0.05); + --card-bg: rgba(255, 255, 255, 0.025); + --ring-bg: rgba(255, 255, 255, 0.04); + --divider: rgba(255, 255, 255, 0.05); + --input-bg: rgba(0, 0, 0, 0.4); + --radius: 18px; + } + + [data-theme="light"] { + --bg-deep: #eef1f5; + --bg-surface: rgba(255, 255, 255, 0.9); + --border-glow: rgba(0, 160, 200, 0.15); + --accent: #0090c0; + --accent2: #6a20e0; + --accent3: #d02060; + --accent4: #00a060; + --text-primary: #1a1d30; + --text-secondary: #5a6088; + --glass: rgba(255, 255, 255, 0.8); + --cell-bg: rgba(0, 0, 0, 0.03); + --cell-border: rgba(0, 0, 0, 0.06); + --card-bg: rgba(0, 0, 0, 0.03); + --ring-bg: rgba(0, 0, 0, 0.06); + --divider: rgba(0, 0, 0, 0.06); + --input-bg: rgba(0, 0, 0, 0.05); + } + + body { + font-family: 'JetBrains Mono', monospace; + background: var(--bg-deep); + color: var(--text-primary); + min-height: 100vh; + overflow-x: hidden; + } + + /* Animated grid background */ + .bg-grid { + position: fixed; inset: 0; z-index: 0; + background-image: + linear-gradient(rgba(0, 240, 255, 0.025) 1px, transparent 1px), + linear-gradient(90deg, rgba(0, 240, 255, 0.025) 1px, transparent 1px); + background-size: 64px 64px; + animation: gridScroll 20s linear infinite; + pointer-events: none; + } + @keyframes gridScroll { + 0% { background-position: 0 0; } + 100% { background-position: 64px 64px; } + } + + /* Floating orbs */ + .orb { + position: fixed; border-radius: 50%; filter: blur(140px); pointer-events: none; + animation: orbFloat 18s ease-in-out infinite; + } + .orb-1 { width: 600px; height: 600px; background: rgba(0, 240, 255, 0.04); top: -250px; left: -150px; } + .orb-2 { width: 500px; height: 500px; background: rgba(123, 47, 255, 0.04); bottom: -200px; right: -100px; animation-delay: -6s; } + .orb-3 { width: 400px; height: 400px; background: rgba(255, 45, 120, 0.03); top: 40%; left: 40%; animation-delay: -12s; } + @keyframes orbFloat { + 0%, 100% { transform: translate(0, 0) scale(1); } + 33% { transform: translate(50px, -40px) scale(1.06); } + 66% { transform: translate(-30px, 30px) scale(0.94); } + } + + .container { + position: relative; z-index: 1; + max-width: 1600px; margin: 0 auto; padding: 24px 20px; + } + + /* Header */ + .header { + display: flex; align-items: center; justify-content: space-between; + padding: 18px 28px; margin-bottom: 24px; + background: var(--glass); + backdrop-filter: blur(24px); -webkit-backdrop-filter: blur(24px); + border: 1px solid var(--border-glow); border-radius: var(--radius); + } + .header-left { display: flex; align-items: center; gap: 14px; } + .logo-icon { + width: 44px; height: 44px; + background: linear-gradient(135deg, var(--accent), var(--accent2)); + border-radius: 11px; font-size: 22px; + display: flex; align-items: center; justify-content: center; + box-shadow: 0 0 28px rgba(0, 240, 255, 0.35); + animation: logoPulse 3s ease-in-out infinite; + } + @keyframes logoPulse { + 0%, 100% { box-shadow: 0 0 22px rgba(0, 240, 255, 0.3); } + 50% { box-shadow: 0 0 44px rgba(0, 240, 255, 0.6); } + } + .header h1 { + font-family: 'Orbitron', sans-serif; font-size: 20px; font-weight: 700; letter-spacing: 3px; + background: linear-gradient(90deg, var(--accent), var(--accent2)); + -webkit-background-clip: text; -webkit-text-fill-color: transparent; background-clip: text; + } + .header-sub { font-size: 10px; color: var(--text-secondary); letter-spacing: 2px; text-transform: uppercase; } + .header-right { display: flex; align-items: center; gap: 18px; } + .header-date { + font-size: 10px; color: var(--text-secondary); letter-spacing: 1px; + text-align: right; + } + .header-date .date-val { + font-family: 'Orbitron', sans-serif; font-size: 14px; color: var(--accent); display: block; + } + .live-chip { + display: flex; align-items: center; gap: 8px; + font-size: 9px; letter-spacing: 1.5px; padding: 5px 12px; border-radius: 12px; + font-weight: 600; text-transform: uppercase; + background: rgba(0, 255, 136, 0.1); color: var(--accent4); + border: 1px solid rgba(0, 255, 136, 0.2); + } + .live-chip .dot { + width: 7px; height: 7px; border-radius: 50%; + background: var(--accent4); box-shadow: 0 0 8px var(--accent4); + animation: blink 1.5s ease-in-out infinite; + } + @keyframes blink { 0%, 100% { opacity: 1; } 50% { opacity: 0.25; } } + + .theme-toggle { + background: none; border: 1px solid var(--cell-border); + border-radius: 8px; padding: 6px 10px; cursor: pointer; + font-size: 16px; line-height: 1; + transition: border-color 0.3s, background 0.3s; + } + .theme-toggle:hover { border-color: var(--accent); background: var(--cell-bg); } + + /* Section title */ + .section-title { + font-family: 'Orbitron', sans-serif; font-size: 13px; font-weight: 600; letter-spacing: 2px; + margin-bottom: 16px; + background: linear-gradient(90deg, var(--accent), var(--accent2)); + -webkit-background-clip: text; -webkit-text-fill-color: transparent; background-clip: text; + display: inline-block; + } + + /* Summary Cards */ + .summary-cards { + display: grid; + grid-template-columns: repeat(4, 1fr); + gap: 16px; + margin-bottom: 24px; + } + .s-card { + background: var(--glass); + backdrop-filter: blur(24px); -webkit-backdrop-filter: blur(24px); + border: 1px solid var(--border-glow); border-radius: var(--radius); + padding: 20px 22px; + transition: border-color 0.3s, transform 0.2s; + position: relative; overflow: hidden; + } + .s-card:hover { border-color: rgba(0, 240, 255, 0.3); transform: translateY(-2px); } + .s-card .s-icon { + width: 42px; height: 42px; border-radius: 12px; + display: flex; align-items: center; justify-content: center; + font-size: 18px; margin-bottom: 14px; + } + .s-card .s-icon.live { background: linear-gradient(135deg, var(--accent4), #059669); color: #fff; } + .s-card .s-icon.primary { background: linear-gradient(135deg, var(--accent), var(--accent2)); color: #000; } + .s-card .s-icon.warm { background: linear-gradient(135deg, var(--accent2), var(--accent3)); color: #fff; } + .s-card .s-icon.success { background: linear-gradient(135deg, var(--accent), var(--accent4)); color: #000; } + .s-card .s-chip { + font-size: 8px; text-transform: uppercase; letter-spacing: 2px; padding: 3px 10px; + border-radius: 10px; font-weight: 600; + position: absolute; top: 16px; right: 16px; + } + .s-card .s-chip.live-chip-tag { + background: rgba(0, 255, 136, 0.1); color: var(--accent4); + border: 1px solid rgba(0, 255, 136, 0.2); + } + .s-card .s-chip.live-chip-tag .ldot { + width: 6px; height: 6px; border-radius: 50%; display: inline-block; + background: var(--accent4); margin-right: 4px; + animation: blink 1.5s ease-in-out infinite; + } + .s-card .s-chip.info { + background: rgba(0, 240, 255, 0.08); color: var(--accent); + border: 1px solid rgba(0, 240, 255, 0.15); + } + .s-card .s-val { + font-family: 'Orbitron', sans-serif; font-size: 32px; font-weight: 700; + background: linear-gradient(180deg, #fff 0%, var(--accent) 100%); + -webkit-background-clip: text; -webkit-text-fill-color: transparent; background-clip: text; + margin-bottom: 4px; + } + .s-card .s-label { font-size: 11px; color: var(--text-secondary); } + .s-card .s-sub { font-size: 9px; color: var(--text-secondary); margin-top: 4px; } + + .hide-me { display: none; } + + /* Charts + Stats layout */ + .charts-row { + display: grid; + grid-template-columns: 2fr 1fr; + gap: 20px; + margin-bottom: 24px; + } + .panel { + background: var(--glass); + backdrop-filter: blur(24px); -webkit-backdrop-filter: blur(24px); + border: 1px solid var(--border-glow); border-radius: var(--radius); + padding: 24px 26px; + position: relative; overflow: hidden; + } + + .chart-header { + display: flex; align-items: center; justify-content: space-between; + margin-bottom: 20px; flex-wrap: wrap; gap: 12px; + } + .chart-filters { display: flex; gap: 6px; } + .chart-filters button { + font-family: 'JetBrains Mono', monospace; font-size: 10px; + letter-spacing: 1px; padding: 6px 14px; border-radius: 7px; + border: 1px solid var(--cell-border); + background: var(--cell-bg); color: var(--text-secondary); + cursor: pointer; transition: all 0.3s; + } + .chart-filters button:hover { border-color: var(--accent); color: var(--text-primary); } + .chart-filters button.active { + background: linear-gradient(135deg, var(--accent), var(--accent2)); + color: #000; border-color: transparent; font-weight: 600; + } + + .chart-container { position: relative; height: 320px; } + + /* Quick Stats */ + .quick-stat { + display: flex; align-items: center; justify-content: space-between; + padding: 14px 16px; margin-bottom: 10px; + background: var(--cell-bg); + border: 1px solid var(--cell-border); border-radius: 10px; + transition: border-color 0.3s; + } + .quick-stat:hover { border-color: var(--border-glow); } + .quick-stat .qs-left { display: flex; align-items: center; gap: 12px; } + .quick-stat .qs-icon { + width: 38px; height: 38px; border-radius: 10px; + display: flex; align-items: center; justify-content: center; + font-size: 15px; + } + .quick-stat .qs-icon.primary { background: linear-gradient(135deg, var(--accent), var(--accent2)); color: #000; } + .quick-stat .qs-icon.warm { background: linear-gradient(135deg, var(--accent2), var(--accent3)); color: #fff; } + .quick-stat .qs-icon.success { background: linear-gradient(135deg, var(--accent), var(--accent4)); color: #000; } + .quick-stat .qs-label { font-size: 10px; color: var(--text-secondary); } + .quick-stat .qs-sub { font-size: 8px; color: var(--text-secondary); } + .quick-stat .qs-val { + font-family: 'Orbitron', sans-serif; font-size: 18px; font-weight: 700; + color: var(--text-primary); + } + + .recent-section { margin-top: 20px; padding-top: 20px; border-top: 1px solid var(--divider); } + .recent-section .rs-title { + font-size: 10px; text-transform: uppercase; letter-spacing: 3px; + color: var(--text-secondary); margin-bottom: 14px; + } + .recent-item { + display: flex; align-items: center; justify-content: space-between; + padding: 10px 14px; margin-bottom: 6px; + background: var(--cell-bg); + border: 1px solid var(--cell-border); border-radius: 9px; + transition: border-color 0.2s; + } + .recent-item:hover { border-color: var(--border-glow); } + .recent-item .ri-left { display: flex; align-items: center; gap: 10px; } + .recent-item .ri-badge { + width: 32px; height: 32px; border-radius: 8px; + display: flex; align-items: center; justify-content: center; + font-family: 'Orbitron', sans-serif; font-size: 10px; font-weight: 700; + background: linear-gradient(135deg, var(--accent), var(--accent2)); + color: #000; + } + .recent-item .ri-date { font-size: 10px; color: var(--text-primary); } + .recent-item .ri-time { font-size: 8px; color: var(--text-secondary); } + .recent-item .ri-count { + font-family: 'Orbitron', sans-serif; font-size: 13px; font-weight: 700; + color: var(--text-primary); + } + .recent-item .ri-dur { font-size: 8px; color: var(--text-secondary); text-align: right; } + + /* Table Panel */ + .table-panel { + background: var(--glass); + backdrop-filter: blur(24px); -webkit-backdrop-filter: blur(24px); + border: 1px solid var(--border-glow); border-radius: var(--radius); + padding: 24px 26px; + margin-bottom: 24px; + } + .table-toolbar { + display: flex; align-items: center; justify-content: space-between; + margin-bottom: 20px; flex-wrap: wrap; gap: 12px; + } + .table-toolbar .toolbar-group { display: flex; align-items: center; gap: 8px; } + + .table-toolbar input[type="date"] { + font-family: 'JetBrains Mono', monospace; font-size: 11px; + background: var(--input-bg); color: var(--text-primary); + border: 1px solid var(--cell-border); border-radius: 7px; + padding: 7px 12px; outline: none; + transition: border-color 0.3s; + } + .table-toolbar input[type="date"]:focus { border-color: var(--accent); } + + .btn { + font-family: 'Orbitron', sans-serif; font-size: 10px; letter-spacing: 2px; + padding: 8px 16px; border-radius: 7px; cursor: pointer; font-weight: 600; + border: none; transition: opacity 0.3s, box-shadow 0.3s; + } + .btn-primary { + background: linear-gradient(135deg, var(--accent), var(--accent2)); + color: #000; + } + .btn-primary:hover { opacity: 0.85; box-shadow: 0 0 20px rgba(0, 240, 255, 0.3); } + .btn-ghost { + background: var(--cell-bg); color: var(--text-secondary); + border: 1px solid var(--cell-border); + } + .btn-ghost:hover { border-color: var(--accent); color: var(--text-primary); } + .btn-export-btn { + background: linear-gradient(135deg, var(--accent), var(--accent4)); + color: #000; + } + .btn-export-btn:hover { opacity: 0.85; box-shadow: 0 0 16px rgba(0, 255, 136, 0.3); } + + table { + width: 100%; border-collapse: collapse; + } + table thead th { + font-size: 8px; text-transform: uppercase; letter-spacing: 2px; + color: var(--text-secondary); text-align: left; + padding-bottom: 12px; border-bottom: 1px solid var(--divider); + } + table tbody td { + padding: 12px 0; font-size: 11px; + border-bottom: 1px solid var(--divider); + color: var(--text-primary); + } + .table-row tr { transition: background 0.2s; cursor: pointer; } + .table-row tr:hover { background: var(--cell-bg); } + + .row-date { display: flex; align-items: center; gap: 10px; } + .row-date .rd-icon { + width: 36px; height: 36px; border-radius: 10px; + display: flex; align-items: center; justify-content: center; + font-size: 13px; + background: var(--cell-bg); color: var(--text-secondary); + border: 1px solid var(--cell-border); + } + .row-date .rd-icon.today { + background: linear-gradient(135deg, var(--accent), var(--accent2)); + color: #000; border: none; + } + .badge { + display: inline-flex; align-items: center; padding: 3px 10px; + border-radius: 10px; font-size: 9px; font-weight: 600; letter-spacing: 1px; + } + .badge-batches { background: rgba(0, 240, 255, 0.08); color: var(--accent); } + .badge-green { background: rgba(0, 255, 136, 0.1); color: var(--accent4); } + .badge-gray { background: var(--cell-bg); color: var(--text-secondary); } + .link-text { font-size: 10px; letter-spacing: 1px; cursor: pointer; transition: color 0.2s; } + .link-text:hover { color: var(--accent2); } + + /* Modal */ + .modal-overlay { + position: fixed; inset: 0; z-index: 100; + background: rgba(0, 0, 0, 0.6); + backdrop-filter: blur(8px); -webkit-backdrop-filter: blur(8px); + display: flex; align-items: center; justify-content: center; + padding: 20px; + } + .modal-overlay.hidden { display: none; } + .modal-content { + background: var(--bg-surface); + backdrop-filter: blur(24px); -webkit-backdrop-filter: blur(24px); + border: 1px solid var(--border-glow); border-radius: var(--radius); + padding: 28px 30px; + width: 100%; max-width: 900px; max-height: 85vh; overflow-y: auto; + transform: scale(0.95); opacity: 0; + transition: transform 0.2s ease-out, opacity 0.2s ease-out; + } + .modal-content.open { transform: scale(1); opacity: 1; } + .modal-header { + display: flex; align-items: center; justify-content: space-between; + margin-bottom: 20px; padding-bottom: 16px; + border-bottom: 1px solid var(--divider); + } + .modal-header h2 { + font-family: 'Orbitron', sans-serif; font-size: 16px; letter-spacing: 2px; + background: linear-gradient(90deg, var(--accent), var(--accent2)); + -webkit-background-clip: text; -webkit-text-fill-color: transparent; background-clip: text; + } + .modal-header .close-btn { + background: none; border: 1px solid var(--cell-border); color: var(--text-secondary); + width: 32px; height: 32px; border-radius: 8px; cursor: pointer; + font-size: 16px; display: flex; align-items: center; justify-content: center; + transition: border-color 0.3s, color 0.3s; + } + .modal-header .close-btn:hover { border-color: var(--accent3); color: var(--accent3); } + .modal-subtitle { font-size: 10px; color: var(--text-secondary); letter-spacing: 1px; margin-top: 4px; } + + .modal-stats { + display: grid; grid-template-columns: repeat(3, 1fr); gap: 12px; + margin-bottom: 20px; + } + .modal-stat { + background: var(--cell-bg); + border: 1px solid var(--cell-border); border-radius: 10px; + padding: 16px; text-align: center; + } + .modal-stat .ms-val { + font-family: 'Orbitron', sans-serif; font-size: 24px; font-weight: 700; + color: var(--text-primary); + } + .modal-stat .ms-label { font-size: 9px; color: var(--text-secondary); letter-spacing: 1px; margin-top: 4px; } + .modal-table-wrap { max-height: 380px; overflow-y: auto; } + + /* Footer */ + .footer { + margin-top: 24px; text-align: center; + font-size: 9px; letter-spacing: 2px; color: var(--text-secondary); + text-transform: uppercase; + } + + /* Scrollbar */ + ::-webkit-scrollbar { width: 6px; } + ::-webkit-scrollbar-track { background: transparent; } + ::-webkit-scrollbar-thumb { background: var(--text-secondary); border-radius: 3px; } + + /* Responsive */ + @media (max-width: 1200px) { + .summary-cards { grid-template-columns: repeat(2, 1fr); } + .charts-row { grid-template-columns: 1fr; } + } + @media (max-width: 768px) { + .container { padding: 16px 12px; } + .summary-cards { grid-template-columns: 1fr; } + .header { flex-direction: column; gap: 10px; padding: 14px 20px; } + .header-left { flex-direction: column; text-align: center; } + .header-right { flex-wrap: wrap; justify-content: center; } + .s-card .s-val { font-size: 24px; } + .panel { padding: 18px 16px; } + .modal-content { padding: 18px 16px; max-width: 95vw; } + .modal-stats { grid-template-columns: 1fr; } + .table-toolbar { flex-direction: column; align-items: flex-start; } + } + @media (max-width: 480px) { + .container { padding: 10px 8px; } + .header { padding: 12px 14px; border-radius: 12px; } + .header h1 { font-size: 16px; letter-spacing: 2px; } + .logo-icon { width: 36px; height: 36px; font-size: 18px; border-radius: 9px; } + .s-card { padding: 16px; border-radius: 14px; } + .s-card .s-val { font-size: 22px; } + .panel { padding: 14px 12px; border-radius: 14px; } + .chart-container { height: 240px; } + .footer { font-size: 8px; letter-spacing: 1px; } + } + - - -
-
-
-
-
- -
-
-

ZenAI APC Dashboard

-

Real-time batch counting analytics

-
-
-
-
- Counting Day -

--

-
-
-
- Live -
-
-
-
-
+ -
- -
- -
-
- - LIVE - -
-
-
- -
- Batch #-- -
-

--

-

Current Batch Count

-
- -- -
-
+
+
+
+
- -
-
-
- -
- Today -
-

--

-

Total Ayam Potong

-
- -- batches -
-
- - -
-
-
- -
- Yesterday -
-

--

-

Total Ayam Potong

-
- -- batches -
-
- - -
-
-
- -
- Average -
-

--

-

Per Day

-
- -- days recorded -
-
- - -
-
-
- -
- Best Day -
-

--

-

--

-
- Record -
-
-
- - -
- -
-
-
-

Daily Trends

-

Object count over time

-
-
- - - -
-
-
- -
-
- - -
-

Quick Stats

-
-
-
-
- -
-
-

Grand Total

-

All time

-
-
- -- -
- -
-
-
- -
-
-

Total Batches

-

All time

-
-
- -- -
- -
-
-
- -
-
-

Avg per Batch

-

Overall average

-
-
- -- -
-
- -
-

Recent Activity

-
- -
-
-
-
- - -
-
-
-

Daily Records

-

Click on a row to view batch details

-
-
- - - - - -
-
- -
- - - - - - - - - - - - - - -
DateTotal CountBatchesAvg/BatchStatusAction
-
-
-
- - -