update untuk multiple source

This commit is contained in:
Alberto-Audrix committed 2026-08-11 15:51:38 +07:00
1 parent 8bbe5950ea
commit cfc8365dfb
5 files changed
+255 -52

No files matched your search

+1
View File
@@ -130,6 +130,7 @@ celerybeat.pid
# Environments # Environments
.env .env
sources.env
.venv .venv
env/ env/
venv/ venv/
+46 -15
View File
@@ -4,15 +4,44 @@ Admin dashboard that aggregates **KTC** (tuang) and **KPC** (masuk) counter data
## Sources ## Sources
| Role | Project | Default URL | Cameras | | Role | Project | Typical port | Cameras |
|------|---------|-------------|---------| |------|---------|--------------|---------|
| Tuang | `zenai-ktc-python` | `http://192.168.192.15:9000` | `K1-L1-left`, `K1-L1-right` | | Tuang | `zenai-ktc-python` | `9000` | `K1-L1-left`, `K1-L1-right` |
| Masuk | `zenai-kpc-python` | `http://192.168.192.15:7000` | `K1-L1-in`, `K1-L1-out` | | 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 ## Install path
Deploy code to `/opt/karung-web-admin` (matches the systemd units below). 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) ## Puller (systemd service)
```bash ```bash
@@ -21,7 +50,7 @@ sudo systemctl daemon-reload
sudo systemctl enable --now karung-web-admin-pull-sources.service 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: Writes:
@@ -62,13 +91,15 @@ Combined API: `GET /api/k1/combined?history_days=30`
## Environment ## Environment
| Variable | Default | Configured in `sources.env` (see `sources.env.example`):
|----------|---------|
| `KTC_BASE_URL` | `http://192.168.192.15:9000` | | Variable | Description |
| `KPC_BASE_URL` | `http://192.168.192.15:7000` | |----------|-------------|
| `SOURCE_HISTORY_DAYS` | `30` | | `KTC_BASE_URL` | Comma-separated tuang/KTC base URLs (summed) |
| `SOURCE_PULL_INTERVAL_SECONDS` | `30` | | `KPC_BASE_URL` | Comma-separated masuk/KPC base URLs (summed) |
| `KARUNG_TUANG_JSON` | `karung_tuang.json` | | `SOURCE_HISTORY_DAYS` | History window to pull (default `30`) |
| `KARUNG_MASUK_JSON` | `karung_masuk.json` | | `SOURCE_PULL_INTERVAL_SECONDS` | Pull loop interval (default `30`) |
| `KARUNG_TUANG_DB` | `karung_tuang.db` | | `KARUNG_TUANG_JSON` | Live tuang JSON path (default `karung_tuang.json`) |
| `KARUNG_MASUK_DB` | `karung_masuk.db` | | `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`) |
+2 -4
View File
@@ -5,10 +5,8 @@ After=network.target
[Service] [Service]
Type=simple Type=simple
WorkingDirectory=/opt/karung-web-admin WorkingDirectory=/opt/karung-web-admin
Environment=KTC_BASE_URL=http://192.168.192.15:9000 # Edit IPs in /opt/karung-web-admin/sources.env (copy from sources.env.example)
Environment=KPC_BASE_URL=http://192.168.192.15:7000 EnvironmentFile=/opt/karung-web-admin/sources.env
Environment=SOURCE_HISTORY_DAYS=30
Environment=SOURCE_PULL_INTERVAL_SECONDS=30
ExecStart=/usr/bin/python3 -u /opt/karung-web-admin/pull_sources.py --loop ExecStart=/usr/bin/python3 -u /opt/karung-web-admin/pull_sources.py --loop
Restart=always Restart=always
RestartSec=5s RestartSec=5s
+190 -33
View File
@@ -1,6 +1,8 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
""" """
Pull counter data from KTC (tuang) and KPC (masuk) dashboards into local JSON + SQLite. 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 from __future__ import annotations
@@ -17,11 +19,10 @@ from pathlib import Path
BASE_DIR = Path(__file__).resolve().parent 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_HISTORY_DAYS = 30
DEFAULT_PULL_INTERVAL_SECONDS = 30 DEFAULT_PULL_INTERVAL_SECONDS = 30
PULL_STATUS_FILE = BASE_DIR / ".pull_status.json" PULL_STATUS_FILE = BASE_DIR / ".pull_status.json"
SOURCES_ENV_FILE = BASE_DIR / "sources.env"
KTC_LIVE_FIELDS = ( KTC_LIVE_FIELDS = (
("count_left", "K1-L1-left"), ("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") 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): def get_configured_path(env_name, default_file_name):
path = Path(os.environ.get(env_name, default_file_name)).expanduser() path = Path(os.environ.get(env_name, default_file_name)).expanduser()
if not path.is_absolute(): if not path.is_absolute():
@@ -59,6 +75,30 @@ def get_int_env(env_name, default_value):
return int(raw) 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): def map_live_payload(payload, field_map):
"""Map /api/current-counter fields to live JSON camera entries.""" """Map /api/current-counter fields to live JSON camera entries."""
mapped = {} mapped = {}
@@ -88,6 +128,37 @@ def map_history_rows(rows, field_map):
return mapped_rows 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): def write_live_json(json_path, data):
json_path = Path(json_path) json_path = Path(json_path)
json_path.parent.mkdir(parents=True, exist_ok=True) json_path.parent.mkdir(parents=True, exist_ok=True)
@@ -204,10 +275,49 @@ def save_pull_status(status):
handle.write("\n") 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( def pull_source(
*, *,
source_key, source_key,
base_url, base_urls,
live_fields, live_fields,
history_fields, history_fields,
json_path, json_path,
@@ -215,38 +325,81 @@ def pull_source(
table_name, table_name,
history_days, history_days,
): ):
result = { """Pull one or more hosts for a role and write summed live/history totals."""
"status": "ok", 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(), "pulled_at": get_utc_timestamp(),
"base_url": base_url.rstrip("/"), "base_url": base_urls[0] if len(base_urls) == 1 else None,
"error": 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( def format_role_status(role_status):
f"{base_url.rstrip('/')}/api/history?days={history_days}" status = role_status.get("status", "unknown")
) hosts_ok = role_status.get("hosts_ok")
history_rows = history_payload.get("data", []) if isinstance(history_payload, dict) else [] hosts_total = role_status.get("hosts_total")
mapped_history = map_history_rows(history_rows, history_fields) detail = ""
if mapped_history: if hosts_ok is not None and hosts_total is not None:
upsert_counter_rows(db_path, table_name, mapped_history) detail = f" {hosts_ok}/{hosts_total} hosts"
error = role_status.get("error")
result["live_cameras"] = list(live_data.keys()) if error:
result["history_rows"] = len(mapped_history) return f"{status}{detail} ({error})"
return result return f"{status}{detail}"
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(): def run_pull():
ktc_base = os.environ.get("KTC_BASE_URL", DEFAULT_KTC_BASE_URL).strip() or DEFAULT_KTC_BASE_URL ktc_urls = get_required_base_urls("KTC_BASE_URL")
kpc_base = os.environ.get("KPC_BASE_URL", DEFAULT_KPC_BASE_URL).strip() or DEFAULT_KPC_BASE_URL kpc_urls = get_required_base_urls("KPC_BASE_URL")
history_days = get_int_env("SOURCE_HISTORY_DAYS", DEFAULT_HISTORY_DAYS) history_days = get_int_env("SOURCE_HISTORY_DAYS", DEFAULT_HISTORY_DAYS)
tuang_json = get_configured_path("KARUNG_TUANG_JSON", "karung_tuang.json") tuang_json = get_configured_path("KARUNG_TUANG_JSON", "karung_tuang.json")
@@ -257,7 +410,7 @@ def run_pull():
status = load_pull_status() status = load_pull_status()
status["ktc"] = pull_source( status["ktc"] = pull_source(
source_key="ktc", source_key="ktc",
base_url=ktc_base, base_urls=ktc_urls,
live_fields=KTC_LIVE_FIELDS, live_fields=KTC_LIVE_FIELDS,
history_fields=KTC_HISTORY_FIELDS, history_fields=KTC_HISTORY_FIELDS,
json_path=tuang_json, json_path=tuang_json,
@@ -267,7 +420,7 @@ def run_pull():
) )
status["kpc"] = pull_source( status["kpc"] = pull_source(
source_key="kpc", source_key="kpc",
base_url=kpc_base, base_urls=kpc_urls,
live_fields=KPC_LIVE_FIELDS, live_fields=KPC_LIVE_FIELDS,
history_fields=KPC_HISTORY_FIELDS, history_fields=KPC_HISTORY_FIELDS,
json_path=masuk_json, json_path=masuk_json,
@@ -278,10 +431,10 @@ def run_pull():
status["updated_at"] = get_utc_timestamp() status["updated_at"] = get_utc_timestamp()
save_pull_status(status) save_pull_status(status)
print(f"KTC: {status['ktc']['status']}" + (f" ({status['ktc']['error']})" if status["ktc"].get("error") else "")) print(f"KTC: {format_role_status(status['ktc'])}")
print(f"KPC: {status['kpc']['status']}" + (f" ({status['kpc']['error']})" if status["kpc"].get("error") else "")) 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 1
return 0 return 0
@@ -299,11 +452,15 @@ def run_pull_loop(interval_seconds=None):
def main(argv=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:]) args = list(argv if argv is not None else sys.argv[1:])
if "--loop" in args: if "--loop" in args:
run_pull_loop() run_pull_loop()
return 0 return 0
return run_pull() return run_pull()
if __name__ == "__main__": if __name__ == "__main__":
raise SystemExit(main()) raise SystemExit(main())
+16
View File
@@ -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