From cfc8365dfba16c79780658bb3a55297ec4fee8ab Mon Sep 17 00:00:00 2001 From: Alberto-Audrix Date: Tue, 11 Aug 2026 15:51:38 +0700 Subject: [PATCH] update untuk multiple source --- .gitignore | 1 + README.md | 61 +++++-- karung-web-admin-pull-sources.service | 6 +- pull_sources.py | 223 ++++++++++++++++++++++---- sources.env.example | 16 ++ 5 files changed, 255 insertions(+), 52 deletions(-) create mode 100644 sources.env.example diff --git a/.gitignore b/.gitignore index 36b13f1..34f50e6 100644 --- a/.gitignore +++ b/.gitignore @@ -130,6 +130,7 @@ celerybeat.pid # Environments .env +sources.env .venv env/ venv/ diff --git a/README.md b/README.md index 416c4ee..a55058a 100644 --- a/README.md +++ b/README.md @@ -4,15 +4,44 @@ Admin dashboard that aggregates **KTC** (tuang) and **KPC** (masuk) counter data ## Sources -| Role | Project | Default URL | Cameras | -|------|---------|-------------|---------| -| Tuang | `zenai-ktc-python` | `http://192.168.192.15:9000` | `K1-L1-left`, `K1-L1-right` | -| Masuk | `zenai-kpc-python` | `http://192.168.192.15:7000` | `K1-L1-in`, `K1-L1-out` | +| Role | Project | Typical port | Cameras | +|------|---------|--------------|---------| +| Tuang | `zenai-ktc-python` | `9000` | `K1-L1-left`, `K1-L1-right` | +| Masuk | `zenai-kpc-python` | `7000` | `K1-L1-in`, `K1-L1-out` | + +Host IPs are **not** hardcoded. Configure them per machine in `sources.env`. + +Multiple hosts per role (comma-separated) are **summed** into one site total. ## Install path Deploy code to `/opt/karung-web-admin` (matches the systemd units below). +## Configure source hosts (required) + +```bash +cd /opt/karung-web-admin +cp sources.env.example sources.env +nano sources.env # set real KTC/KPC URLs for this site +``` + +Example `sources.env`: + +```bash +KTC_BASE_URL=http://10.0.0.11:9000,http://10.0.0.12:9000 +KPC_BASE_URL=http://10.0.0.11:7000,http://10.0.0.12:7000 +SOURCE_HISTORY_DAYS=30 +SOURCE_PULL_INTERVAL_SECONDS=30 +``` + +After editing: + +```bash +sudo systemctl restart karung-web-admin-pull-sources.service +``` + +`sources.env` is gitignored so each production host keeps its own IPs. + ## Puller (systemd service) ```bash @@ -21,7 +50,7 @@ sudo systemctl daemon-reload sudo systemctl enable --now karung-web-admin-pull-sources.service ``` -It runs `pull_sources.py --loop` and refreshes every `SOURCE_PULL_INTERVAL_SECONDS` (default 30). +It loads `/opt/karung-web-admin/sources.env`, runs `pull_sources.py --loop`, and refreshes every `SOURCE_PULL_INTERVAL_SECONDS` (default 30). Writes: @@ -62,13 +91,15 @@ Combined API: `GET /api/k1/combined?history_days=30` ## Environment -| Variable | Default | -|----------|---------| -| `KTC_BASE_URL` | `http://192.168.192.15:9000` | -| `KPC_BASE_URL` | `http://192.168.192.15:7000` | -| `SOURCE_HISTORY_DAYS` | `30` | -| `SOURCE_PULL_INTERVAL_SECONDS` | `30` | -| `KARUNG_TUANG_JSON` | `karung_tuang.json` | -| `KARUNG_MASUK_JSON` | `karung_masuk.json` | -| `KARUNG_TUANG_DB` | `karung_tuang.db` | -| `KARUNG_MASUK_DB` | `karung_masuk.db` | +Configured in `sources.env` (see `sources.env.example`): + +| Variable | Description | +|----------|-------------| +| `KTC_BASE_URL` | Comma-separated tuang/KTC base URLs (summed) | +| `KPC_BASE_URL` | Comma-separated masuk/KPC base URLs (summed) | +| `SOURCE_HISTORY_DAYS` | History window to pull (default `30`) | +| `SOURCE_PULL_INTERVAL_SECONDS` | Pull loop interval (default `30`) | +| `KARUNG_TUANG_JSON` | Live tuang JSON path (default `karung_tuang.json`) | +| `KARUNG_MASUK_JSON` | Live masuk JSON path (default `karung_masuk.json`) | +| `KARUNG_TUANG_DB` | Tuang history DB path (default `karung_tuang.db`) | +| `KARUNG_MASUK_DB` | Masuk history DB path (default `karung_masuk.db`) | diff --git a/karung-web-admin-pull-sources.service b/karung-web-admin-pull-sources.service index da60d5e..0852073 100644 --- a/karung-web-admin-pull-sources.service +++ b/karung-web-admin-pull-sources.service @@ -5,10 +5,8 @@ After=network.target [Service] Type=simple WorkingDirectory=/opt/karung-web-admin -Environment=KTC_BASE_URL=http://192.168.192.15:9000 -Environment=KPC_BASE_URL=http://192.168.192.15:7000 -Environment=SOURCE_HISTORY_DAYS=30 -Environment=SOURCE_PULL_INTERVAL_SECONDS=30 +# Edit IPs in /opt/karung-web-admin/sources.env (copy from sources.env.example) +EnvironmentFile=/opt/karung-web-admin/sources.env ExecStart=/usr/bin/python3 -u /opt/karung-web-admin/pull_sources.py --loop Restart=always RestartSec=5s diff --git a/pull_sources.py b/pull_sources.py index 4097ea2..dc0409f 100644 --- a/pull_sources.py +++ b/pull_sources.py @@ -1,6 +1,8 @@ #!/usr/bin/env python3 """ Pull counter data from KTC (tuang) and KPC (masuk) dashboards into local JSON + SQLite. + +Multiple base URLs (comma-separated) are summed into one site total per role. """ from __future__ import annotations @@ -17,11 +19,10 @@ 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" +SOURCES_ENV_FILE = BASE_DIR / "sources.env" KTC_LIVE_FIELDS = ( ("count_left", "K1-L1-left"), @@ -45,6 +46,21 @@ 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(): @@ -59,6 +75,30 @@ def get_int_env(env_name, default_value): return int(raw) +def parse_base_urls(raw_value): + """Parse comma-separated base URLs.""" + urls = [] + for part in str(raw_value or "").split(","): + url = part.strip().rstrip("/") + if url: + urls.append(url) + return urls + + +def get_required_base_urls(env_name): + """Read required comma-separated source URLs from the environment.""" + urls = parse_base_urls(os.environ.get(env_name, "")) + if urls: + return urls + + 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 map_live_payload(payload, field_map): """Map /api/current-counter fields to live JSON camera entries.""" mapped = {} @@ -88,6 +128,37 @@ def map_history_rows(rows, field_map): return mapped_rows +def sum_live_data(live_payloads): + """Sum karung counts by camera 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) @@ -204,10 +275,49 @@ def save_pull_status(status): handle.write("\n") +def fetch_host_payload(*, base_url, live_fields, history_fields, history_days): + """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) + + 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) + + host_result["live_data"] = live_data + host_result["history_rows"] = mapped_history + host_result["live_cameras"] = 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, - base_url, + base_urls, live_fields, history_fields, json_path, @@ -215,38 +325,81 @@ def pull_source( table_name, history_days, ): - result = { - "status": "ok", + """Pull one or more hosts for a role and write summed live/history totals.""" + hosts = [] + live_payloads = [] + history_groups = [] + + for base_url in base_urls: + host = fetch_host_payload( + base_url=base_url, + live_fields=live_fields, + history_fields=history_fields, + history_days=history_days, + ) + hosts.append( + { + "status": host["status"], + "pulled_at": host["pulled_at"], + "base_url": host["base_url"], + "error": host["error"], + "live_cameras": host.get("live_cameras", []), + "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_url.rstrip("/"), - "error": None, + "base_url": base_urls[0] if len(base_urls) == 1 else None, + "base_urls": list(base_urls), + "combine": "sum", + "hosts": hosts, + "error": "; ".join(errors) if errors else None, + "live_cameras": list(combined_live.keys()), + "history_rows": len(combined_history), + "hosts_ok": ok_count, + "hosts_total": len(hosts), + "source_key": source_key, } - 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 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_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 + ktc_urls = get_required_base_urls("KTC_BASE_URL") + kpc_urls = get_required_base_urls("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") @@ -257,7 +410,7 @@ def run_pull(): status = load_pull_status() status["ktc"] = pull_source( source_key="ktc", - base_url=ktc_base, + base_urls=ktc_urls, live_fields=KTC_LIVE_FIELDS, history_fields=KTC_HISTORY_FIELDS, json_path=tuang_json, @@ -267,7 +420,7 @@ def run_pull(): ) status["kpc"] = pull_source( source_key="kpc", - base_url=kpc_base, + base_urls=kpc_urls, live_fields=KPC_LIVE_FIELDS, history_fields=KPC_HISTORY_FIELDS, json_path=masuk_json, @@ -278,10 +431,10 @@ def run_pull(): 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 "")) + print(f"KTC: {format_role_status(status['ktc'])}") + print(f"KPC: {format_role_status(status['kpc'])}") - if status["ktc"]["status"] != "ok" or status["kpc"]["status"] != "ok": + if status["ktc"]["status"] == "error" or status["kpc"]["status"] == "error": return 1 return 0 @@ -299,11 +452,15 @@ def run_pull_loop(interval_seconds=None): 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()) diff --git a/sources.env.example b/sources.env.example new file mode 100644 index 0000000..a064785 --- /dev/null +++ b/sources.env.example @@ -0,0 +1,16 @@ +# Site source URLs for karung-web-admin +# Copy this file once per machine, then edit the IPs for that site: +# cp sources.env.example sources.env +# +# Comma-separated hosts are summed into one site total. +# Restart after changes: +# sudo systemctl restart karung-web-admin-pull-sources.service + +# Tuang (KTC) +KTC_BASE_URL=http://REPLACE_HOST_A:9000,http://REPLACE_HOST_B:9000 + +# Masuk (KPC) +KPC_BASE_URL=http://REPLACE_HOST_A:7000,http://REPLACE_HOST_B:7000 + +SOURCE_HISTORY_DAYS=30 +SOURCE_PULL_INTERVAL_SECONDS=30