#!/usr/bin/env python3 """ Pull counter data from KTC (tuang) and KPC (masuk) dashboards into local JSON + SQLite. """ from __future__ import annotations import json import os import sqlite3 import sys import time import urllib.error import urllib.request from datetime import datetime, timezone from pathlib import Path BASE_DIR = Path(__file__).resolve().parent DEFAULT_KTC_BASE_URL = "http://192.168.192.15:9000" DEFAULT_KPC_BASE_URL = "http://192.168.192.15:7000" DEFAULT_HISTORY_DAYS = 30 DEFAULT_PULL_INTERVAL_SECONDS = 30 PULL_STATUS_FILE = BASE_DIR / ".pull_status.json" KTC_LIVE_FIELDS = ( ("count_left", "K1-L1-left"), ("count_right", "K1-L1-right"), ) KPC_LIVE_FIELDS = ( ("count_in", "K1-L1-in"), ("count_out", "K1-L1-out"), ) KTC_HISTORY_FIELDS = ( ("total_left", "K1-L1-left"), ("total_right", "K1-L1-right"), ) KPC_HISTORY_FIELDS = ( ("total_in", "K1-L1-in"), ("total_out", "K1-L1-out"), ) def get_utc_timestamp(): return datetime.now(timezone.utc).isoformat(timespec="seconds").replace("+00:00", "Z") def get_configured_path(env_name, default_file_name): path = Path(os.environ.get(env_name, default_file_name)).expanduser() if not path.is_absolute(): path = BASE_DIR / path return path def get_int_env(env_name, default_value): raw = os.environ.get(env_name) if raw is None or raw.strip() == "": return default_value return int(raw) def map_live_payload(payload, field_map): """Map /api/current-counter fields to live JSON camera entries.""" mapped = {} for field_name, camera_name in field_map: value = payload.get(field_name, 0) try: karung = int(value or 0) except (TypeError, ValueError): karung = 0 mapped[camera_name] = {"karung": karung} return mapped def map_history_rows(rows, field_map): """Map /api/history data rows to (camera_name, date, counter_value) tuples.""" mapped_rows = [] for row in rows or []: counting_date = row.get("date") or row.get("counting_date") if not counting_date: continue for field_name, camera_name in field_map: try: value = int(row.get(field_name, 0) or 0) except (TypeError, ValueError): value = 0 mapped_rows.append((camera_name, str(counting_date), value)) return mapped_rows def write_live_json(json_path, data): json_path = Path(json_path) json_path.parent.mkdir(parents=True, exist_ok=True) with json_path.open("w", encoding="utf-8") as handle: json.dump(data, handle, indent=2) handle.write("\n") def ensure_tuang_table(conn): conn.execute( """ CREATE TABLE IF NOT EXISTS counter_data ( id INTEGER PRIMARY KEY AUTOINCREMENT, camera_name TEXT NOT NULL, date TIMESTAMP NOT NULL, counter_value INTEGER NOT NULL ) """ ) def ensure_masuk_table(conn): conn.execute( """ CREATE TABLE IF NOT EXISTS karung_counts ( id INTEGER PRIMARY KEY AUTOINCREMENT, camera_name TEXT NOT NULL, date DATE NOT NULL, counter_value INTEGER NOT NULL, timestamp DATETIME DEFAULT CURRENT_TIMESTAMP ) """ ) def upsert_counter_rows(db_path, table_name, rows): """Upsert counter rows by (camera_name, date).""" db_path = Path(db_path) db_path.parent.mkdir(parents=True, exist_ok=True) conn = sqlite3.connect(db_path) try: if table_name == "counter_data": ensure_tuang_table(conn) date_expr = "date(date)" else: ensure_masuk_table(conn) date_expr = "date" for camera_name, counting_date, counter_value in rows: existing = conn.execute( f""" SELECT id FROM {table_name} WHERE camera_name = ? AND {date_expr} = ? LIMIT 1 """, (camera_name, counting_date), ).fetchone() if existing: conn.execute( f""" UPDATE {table_name} SET counter_value = ? WHERE id = ? """, (counter_value, existing[0]), ) else: if table_name == "karung_counts": conn.execute( """ INSERT INTO karung_counts (camera_name, date, counter_value) VALUES (?, ?, ?) """, (camera_name, counting_date, counter_value), ) else: conn.execute( """ INSERT INTO counter_data (camera_name, date, counter_value) VALUES (?, ?, ?) """, (camera_name, counting_date, counter_value), ) conn.commit() finally: conn.close() def fetch_json(url, timeout=30): request = urllib.request.Request(url, method="GET") request.add_header("Accept", "application/json") with urllib.request.urlopen(request, timeout=timeout) as response: return json.loads(response.read().decode("utf-8")) def load_pull_status(): if not PULL_STATUS_FILE.exists(): return {"ktc": {"status": "unknown"}, "kpc": {"status": "unknown"}} try: with PULL_STATUS_FILE.open(encoding="utf-8") as handle: data = json.load(handle) if isinstance(data, dict): return data except (OSError, json.JSONDecodeError): pass return {"ktc": {"status": "unknown"}, "kpc": {"status": "unknown"}} def save_pull_status(status): with PULL_STATUS_FILE.open("w", encoding="utf-8") as handle: json.dump(status, handle, indent=2) handle.write("\n") def pull_source( *, source_key, base_url, live_fields, history_fields, json_path, db_path, table_name, history_days, ): result = { "status": "ok", "pulled_at": get_utc_timestamp(), "base_url": base_url.rstrip("/"), "error": None, } try: current = fetch_json(f"{base_url.rstrip('/')}/api/current-counter") live_data = map_live_payload(current, live_fields) write_live_json(json_path, live_data) history_payload = fetch_json( f"{base_url.rstrip('/')}/api/history?days={history_days}" ) history_rows = history_payload.get("data", []) if isinstance(history_payload, dict) else [] mapped_history = map_history_rows(history_rows, history_fields) if mapped_history: upsert_counter_rows(db_path, table_name, mapped_history) result["live_cameras"] = list(live_data.keys()) result["history_rows"] = len(mapped_history) return result except (urllib.error.URLError, urllib.error.HTTPError, TimeoutError, json.JSONDecodeError, OSError, ValueError) as error: result["status"] = "error" result["error"] = str(error) return result def run_pull(): ktc_base = os.environ.get("KTC_BASE_URL", DEFAULT_KTC_BASE_URL).strip() or DEFAULT_KTC_BASE_URL kpc_base = os.environ.get("KPC_BASE_URL", DEFAULT_KPC_BASE_URL).strip() or DEFAULT_KPC_BASE_URL history_days = get_int_env("SOURCE_HISTORY_DAYS", DEFAULT_HISTORY_DAYS) tuang_json = get_configured_path("KARUNG_TUANG_JSON", "karung_tuang.json") masuk_json = get_configured_path("KARUNG_MASUK_JSON", "karung_masuk.json") tuang_db = get_configured_path("KARUNG_TUANG_DB", "karung_tuang.db") masuk_db = get_configured_path("KARUNG_MASUK_DB", "karung_masuk.db") status = load_pull_status() status["ktc"] = pull_source( source_key="ktc", base_url=ktc_base, live_fields=KTC_LIVE_FIELDS, history_fields=KTC_HISTORY_FIELDS, json_path=tuang_json, db_path=tuang_db, table_name="counter_data", history_days=history_days, ) status["kpc"] = pull_source( source_key="kpc", base_url=kpc_base, live_fields=KPC_LIVE_FIELDS, history_fields=KPC_HISTORY_FIELDS, json_path=masuk_json, db_path=masuk_db, table_name="karung_counts", history_days=history_days, ) status["updated_at"] = get_utc_timestamp() save_pull_status(status) print(f"KTC: {status['ktc']['status']}" + (f" ({status['ktc']['error']})" if status["ktc"].get("error") else "")) print(f"KPC: {status['kpc']['status']}" + (f" ({status['kpc']['error']})" if status["kpc"].get("error") else "")) if status["ktc"]["status"] != "ok" or status["kpc"]["status"] != "ok": return 1 return 0 def run_pull_loop(interval_seconds=None): if interval_seconds is None: interval_seconds = get_int_env("SOURCE_PULL_INTERVAL_SECONDS", DEFAULT_PULL_INTERVAL_SECONDS) if interval_seconds < 1: interval_seconds = DEFAULT_PULL_INTERVAL_SECONDS print(f"Starting source pull loop every {interval_seconds}s") while True: run_pull() time.sleep(interval_seconds) def main(argv=None): args = list(argv if argv is not None else sys.argv[1:]) if "--loop" in args: run_pull_loop() return 0 return run_pull() if __name__ == "__main__": raise SystemExit(main())