Files
karung-web-admin/pull_sources.py
T

537 lines
17 KiB
Python

#!/usr/bin/env python3
"""
Pull counter data from KTC (tuang) and KPC (masuk) dashboards into local JSON + SQLite.
Source URL groups are labeled by env order:
tuang -> K1 - Use, K2 - Use, ... (API total_count; hosts in a group are summed)
masuk -> K1 - In / K1 - Out, K2 - In / K2 - Out, ...
Grouping syntax in KTC_BASE_URL / KPC_BASE_URL:
- `|` separates cameras (K1, K2, ...)
- `,` separates hosts inside one camera (summed into that Kn)
- No `|` keeps legacy behavior: each comma-separated host is its own Kn
- Single multi-host camera: `http://a,http://b,http://c|` (trailing `|`)
"""
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_HISTORY_DAYS = 30
DEFAULT_PULL_INTERVAL_SECONDS = 30
PULL_STATUS_FILE = BASE_DIR / ".pull_status.json"
SOURCES_ENV_FILE = BASE_DIR / "sources.env"
# Field suffix is combined with source index into names like "K1 - Use".
# KTC uses total_count so hosts without left/right still report correctly.
KTC_LIVE_FIELDS = (
("total_count", "Use"),
)
KPC_LIVE_FIELDS = (
("count_in", "In"),
("count_out", "Out"),
)
KTC_HISTORY_FIELDS = (
("total_count", "Use"),
)
KPC_HISTORY_FIELDS = (
("total_in", "In"),
("total_out", "Out"),
)
def get_utc_timestamp():
return datetime.now(timezone.utc).isoformat(timespec="seconds").replace("+00:00", "Z")
def load_env_file(env_path):
"""Load KEY=VALUE pairs into the process environment (does not override existing)."""
env_path = Path(env_path)
if not env_path.exists():
return
with env_path.open(encoding="utf-8") as env_file:
for line in env_file:
stripped_line = line.strip()
if not stripped_line or stripped_line.startswith("#") or "=" not in stripped_line:
continue
key, value = stripped_line.split("=", 1)
os.environ.setdefault(key.strip(), value.strip().strip('"').strip("'"))
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 parse_base_urls(raw_value):
"""Parse comma-separated base URLs into a flat list."""
urls = []
for part in str(raw_value or "").split(","):
url = part.strip().rstrip("/")
if url:
urls.append(url)
return urls
def parse_source_groups(raw_value):
"""Parse source URL groups for Kn labeling.
With `|`: each segment is one Kn; commas inside a segment are hosts summed into that Kn.
Without `|`: legacy mode — each comma-separated host is its own Kn.
Trailing `|` enables group mode for a single multi-host camera.
"""
raw = str(raw_value or "").strip()
if not raw:
return []
if "|" in raw:
groups = []
for segment in raw.split("|"):
urls = parse_base_urls(segment)
if urls:
groups.append(urls)
return groups
return [[url] for url in parse_base_urls(raw)]
def get_required_source_groups(env_name):
"""Read required source URL groups from the environment."""
groups = parse_source_groups(os.environ.get(env_name, ""))
if groups:
return groups
example_path = BASE_DIR / "sources.env.example"
raise SystemExit(
f"Missing or empty {env_name}. "
f"Copy {example_path.name} to {SOURCES_ENV_FILE.name}, set the host URLs for this site, "
"then restart the puller service."
)
def flatten_source_groups(groups):
"""Flatten groups into a host list (order preserved)."""
return [url for group in groups for url in group]
def build_source_camera_name(source_index, camera_suffix):
return f"K{source_index} - {camera_suffix}"
def extract_history_calendar_date(row):
"""Use the actual activity date from KPC history, not the forward-looking counting_date."""
start_time = row.get("start_time")
if start_time:
try:
parsed = datetime.fromisoformat(str(start_time).replace("Z", "+00:00"))
if parsed.tzinfo is not None:
parsed = parsed.astimezone()
return parsed.date().isoformat()
except ValueError:
pass
counting_date = row.get("date") or row.get("counting_date")
return str(counting_date) if counting_date else None
def map_live_payload(payload, field_map, source_index):
"""Map /api/current-counter fields to live JSON source entries."""
mapped = {}
for field_name, camera_suffix in field_map:
value = payload.get(field_name, 0)
try:
karung = int(value or 0)
except (TypeError, ValueError):
karung = 0
camera_name = build_source_camera_name(source_index, camera_suffix)
entry = mapped.setdefault(camera_name, {"karung": 0})
entry["karung"] += karung
return mapped
def map_history_rows(rows, field_map, source_index):
"""Map /api/history data rows to (camera_name, date, counter_value) tuples."""
mapped_rows = []
for row in rows or []:
counting_date = extract_history_calendar_date(row)
if not counting_date:
continue
totals_by_camera = {}
for field_name, camera_suffix in field_map:
try:
value = int(row.get(field_name, 0) or 0)
except (TypeError, ValueError):
value = 0
camera_name = build_source_camera_name(source_index, camera_suffix)
totals_by_camera[camera_name] = totals_by_camera.get(camera_name, 0) + value
for camera_name, counter_value in totals_by_camera.items():
mapped_rows.append((camera_name, str(counting_date), counter_value))
return mapped_rows
def sum_live_data(live_payloads):
"""Merge live karung counts by source entry name across hosts."""
combined = {}
for live_data in live_payloads:
for camera_name, camera_data in live_data.items():
karung = 0
if isinstance(camera_data, dict):
try:
karung = int(camera_data.get("karung", 0) or 0)
except (TypeError, ValueError):
karung = 0
entry = combined.setdefault(camera_name, {"karung": 0})
entry["karung"] += karung
return combined
def sum_history_rows(history_row_groups):
"""Sum history rows by (camera_name, date) across hosts."""
totals = {}
for rows in history_row_groups:
for camera_name, counting_date, counter_value in rows:
key = (camera_name, counting_date)
totals[key] = totals.get(key, 0) + int(counter_value or 0)
return [
(camera_name, counting_date, value)
for (camera_name, counting_date), value in sorted(
totals.items(), key=lambda item: (item[0][1], item[0][0])
)
]
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 fetch_host_payload(*, base_url, live_fields, history_fields, history_days, source_index):
"""Fetch live + history from one host. Does not write local files."""
host_result = {
"status": "ok",
"pulled_at": get_utc_timestamp(),
"base_url": base_url.rstrip("/"),
"error": None,
"live_data": {},
"history_rows": [],
}
try:
current = fetch_json(f"{base_url.rstrip('/')}/api/current-counter")
live_data = map_live_payload(current, live_fields, source_index)
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, source_index)
host_result["live_data"] = live_data
host_result["history_rows"] = mapped_history
host_result["live_sources"] = list(live_data.keys())
host_result["history_row_count"] = len(mapped_history)
return host_result
except (
urllib.error.URLError,
urllib.error.HTTPError,
TimeoutError,
json.JSONDecodeError,
OSError,
ValueError,
) as error:
host_result["status"] = "error"
host_result["error"] = str(error)
return host_result
def pull_source(
*,
source_key,
source_groups,
live_fields,
history_fields,
json_path,
db_path,
table_name,
history_days,
):
"""Pull grouped hosts for a role and write labeled live/history totals.
Each group shares one Kn label; host totals inside a group are summed.
"""
hosts = []
live_payloads = []
history_groups = []
base_urls = flatten_source_groups(source_groups)
for source_index, group_urls in enumerate(source_groups, start=1):
for base_url in group_urls:
host = fetch_host_payload(
base_url=base_url,
live_fields=live_fields,
history_fields=history_fields,
history_days=history_days,
source_index=source_index,
)
hosts.append(
{
"status": host["status"],
"pulled_at": host["pulled_at"],
"base_url": host["base_url"],
"source_index": source_index,
"error": host["error"],
"live_sources": host.get("live_sources", []),
"history_rows": host.get("history_row_count", 0),
}
)
if host["status"] == "ok":
live_payloads.append(host["live_data"])
history_groups.append(host["history_rows"])
ok_count = sum(1 for host in hosts if host["status"] == "ok")
if ok_count == 0:
overall_status = "error"
elif ok_count < len(hosts):
overall_status = "partial"
else:
overall_status = "ok"
combined_live = sum_live_data(live_payloads)
combined_history = sum_history_rows(history_groups)
if ok_count > 0:
write_live_json(json_path, combined_live)
if combined_history:
upsert_counter_rows(db_path, table_name, combined_history)
errors = [host["error"] for host in hosts if host.get("error")]
return {
"status": overall_status,
"pulled_at": get_utc_timestamp(),
"base_url": base_urls[0] if len(base_urls) == 1 else None,
"base_urls": list(base_urls),
"source_groups": [list(group) for group in source_groups],
"combine": "sum",
"hosts": hosts,
"error": "; ".join(errors) if errors else None,
"live_sources": list(combined_live.keys()),
"history_rows": len(combined_history),
"hosts_ok": ok_count,
"hosts_total": len(hosts),
"groups_total": len(source_groups),
"source_key": source_key,
}
def format_role_status(role_status):
status = role_status.get("status", "unknown")
hosts_ok = role_status.get("hosts_ok")
hosts_total = role_status.get("hosts_total")
detail = ""
if hosts_ok is not None and hosts_total is not None:
detail = f" {hosts_ok}/{hosts_total} hosts"
error = role_status.get("error")
if error:
return f"{status}{detail} ({error})"
return f"{status}{detail}"
def run_pull():
ktc_groups = get_required_source_groups("KTC_BASE_URL")
kpc_groups = get_required_source_groups("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",
source_groups=ktc_groups,
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",
source_groups=kpc_groups,
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: {format_role_status(status['ktc'])}")
print(f"KPC: {format_role_status(status['kpc'])}")
if status["ktc"]["status"] == "error" or status["kpc"]["status"] == "error":
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):
load_env_file(SOURCES_ENV_FILE)
load_env_file(BASE_DIR / ".env")
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())