Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7d9e62f0d7 | ||
|
|
513fe2870f | ||
|
|
014747fd36 | ||
|
|
bb2c358157 | ||
|
|
5b5d28ce52 | ||
|
|
c608580f16 | ||
|
|
fc617b9cf0 | ||
|
|
3823b511ba | ||
|
|
1157d4e8be | ||
|
|
cfc8365dfb | ||
|
|
8bbe5950ea |
No files matched your search
@@ -0,0 +1,64 @@
|
||||
# karung-web-admin environment configuration
|
||||
#
|
||||
# Copy once per machine, then edit values for that site:
|
||||
# cp .env.example .env
|
||||
#
|
||||
# Restart services after changes:
|
||||
# sudo systemctl restart karung-web-admin.service
|
||||
# sudo systemctl restart karung-web-admin-pull-sources.service
|
||||
|
||||
# -----------------------------------------------------------------------------
|
||||
# Source URLs (required for pull_sources.py)
|
||||
# -----------------------------------------------------------------------------
|
||||
# Grouping:
|
||||
# `|` separates cameras (K1, K2, ...)
|
||||
# `,` separates hosts inside one camera (counts are summed into that Kn)
|
||||
# No `|` keeps legacy behavior: each comma-separated host is its own Kn
|
||||
# One multi-host camera only: add a trailing `|` (e.g. http://a,http://b|)
|
||||
|
||||
# Tuang (KTC) — cameras labeled K1 - Use, K2 - Use, ...
|
||||
# Example: K1 = three hosts summed; K2 = one host
|
||||
# KTC_BASE_URL=http://REPLACE_HOST_A:9000,http://REPLACE_HOST_B:9000,http://REPLACE_HOST_C:9000|http://REPLACE_HOST_D:9000
|
||||
KTC_BASE_URL=http://REPLACE_HOST_A:9000|http://REPLACE_HOST_B:9000
|
||||
|
||||
# Masuk (KPC) — cameras labeled K1 - In / K1 - Out, K2 - In / K2 - Out, ...
|
||||
# KPC_BASE_URL=http://REPLACE_HOST_A:7000,http://REPLACE_HOST_B:7000|http://REPLACE_HOST_C:7000
|
||||
KPC_BASE_URL=http://REPLACE_HOST_A:7000|http://REPLACE_HOST_B:7000
|
||||
|
||||
# How many days of history to pull from source APIs (default: 30)
|
||||
SOURCE_HISTORY_DAYS=30
|
||||
|
||||
# Pull loop interval in seconds (default: 30)
|
||||
SOURCE_PULL_INTERVAL_SECONDS=30
|
||||
|
||||
# -----------------------------------------------------------------------------
|
||||
# Data file paths (optional — defaults are in the app directory)
|
||||
# -----------------------------------------------------------------------------
|
||||
# KARUNG_TUANG_JSON=karung_tuang.json
|
||||
# KARUNG_MASUK_JSON=karung_masuk.json
|
||||
# KARUNG_TUANG_DB=karung_tuang.db
|
||||
# KARUNG_MASUK_DB=karung_masuk.db
|
||||
|
||||
# -----------------------------------------------------------------------------
|
||||
# Daily data business-day cutoffs (optional — overrides dashboard DB settings)
|
||||
# -----------------------------------------------------------------------------
|
||||
# Times use HH:MM or HH:MM:SS in LOCAL_TIMEZONE.
|
||||
# Use 24:00:00 for end-of-day (calendar day until midnight).
|
||||
# After cutoff, the dashboard labels that counter side as the next business date.
|
||||
#
|
||||
# LOCAL_TIMEZONE: Asia/Jakarta (WIB), Asia/Makassar (WITA), or Asia/Jayapura (WIT)
|
||||
LOCAL_TIMEZONE=Asia/Jakarta
|
||||
|
||||
# When Karung Tuang (Use) rolls to the next business day
|
||||
TUANG_CUTOFF_TIME=17:00:00
|
||||
|
||||
# When Karung Masuk (In/Out) rolls to the next business day
|
||||
MASUK_CUTOFF_TIME=24:00:00
|
||||
|
||||
# Cycle API (maps K1 -> Kandang 1, K2 -> Kandang 2, etc.)
|
||||
# Base URL without trailing slash; http:// is added automatically if omitted.
|
||||
CYCLE_API_BASE_URL=http://192.168.192.106:5002
|
||||
CYCLE_API_KEY=test-api-key
|
||||
|
||||
# Optional path for exported counter settings JSON (used by other services)
|
||||
# COUNTER_SETTINGS_FILE=counter_settings.json
|
||||
@@ -130,6 +130,7 @@ celerybeat.pid
|
||||
|
||||
# Environments
|
||||
.env
|
||||
sources.env
|
||||
.venv
|
||||
env/
|
||||
venv/
|
||||
|
||||
@@ -0,0 +1,111 @@
|
||||
{
|
||||
"ktc": {
|
||||
"status": "ok",
|
||||
"pulled_at": "2026-09-02T01:00:03Z",
|
||||
"base_url": null,
|
||||
"base_urls": [
|
||||
"http://192.168.192.17:9000",
|
||||
"http://192.168.192.15:9000"
|
||||
],
|
||||
"source_groups": [
|
||||
[
|
||||
"http://192.168.192.17:9000"
|
||||
],
|
||||
[
|
||||
"http://192.168.192.15:9000"
|
||||
]
|
||||
],
|
||||
"combine": "sum",
|
||||
"hosts": [
|
||||
{
|
||||
"status": "ok",
|
||||
"pulled_at": "2026-09-02T01:00:02Z",
|
||||
"base_url": "http://192.168.192.17:9000",
|
||||
"source_index": 1,
|
||||
"error": null,
|
||||
"live_sources": [
|
||||
"K1 - Use"
|
||||
],
|
||||
"history_rows": 2
|
||||
},
|
||||
{
|
||||
"status": "ok",
|
||||
"pulled_at": "2026-09-02T01:00:02Z",
|
||||
"base_url": "http://192.168.192.15:9000",
|
||||
"source_index": 2,
|
||||
"error": null,
|
||||
"live_sources": [
|
||||
"K2 - Use"
|
||||
],
|
||||
"history_rows": 2
|
||||
}
|
||||
],
|
||||
"error": null,
|
||||
"live_sources": [
|
||||
"K1 - Use",
|
||||
"K2 - Use"
|
||||
],
|
||||
"history_rows": 4,
|
||||
"hosts_ok": 2,
|
||||
"hosts_total": 2,
|
||||
"groups_total": 2,
|
||||
"source_key": "ktc"
|
||||
},
|
||||
"kpc": {
|
||||
"status": "ok",
|
||||
"pulled_at": "2026-09-02T01:00:03Z",
|
||||
"base_url": null,
|
||||
"base_urls": [
|
||||
"http://192.168.192.17:7000",
|
||||
"http://192.168.192.15:7000"
|
||||
],
|
||||
"source_groups": [
|
||||
[
|
||||
"http://192.168.192.17:7000"
|
||||
],
|
||||
[
|
||||
"http://192.168.192.15:7000"
|
||||
]
|
||||
],
|
||||
"combine": "sum",
|
||||
"hosts": [
|
||||
{
|
||||
"status": "ok",
|
||||
"pulled_at": "2026-09-02T01:00:03Z",
|
||||
"base_url": "http://192.168.192.17:7000",
|
||||
"source_index": 1,
|
||||
"error": null,
|
||||
"live_sources": [
|
||||
"K1 - In",
|
||||
"K1 - Out"
|
||||
],
|
||||
"history_rows": 6
|
||||
},
|
||||
{
|
||||
"status": "ok",
|
||||
"pulled_at": "2026-09-02T01:00:03Z",
|
||||
"base_url": "http://192.168.192.15:7000",
|
||||
"source_index": 2,
|
||||
"error": null,
|
||||
"live_sources": [
|
||||
"K2 - In",
|
||||
"K2 - Out"
|
||||
],
|
||||
"history_rows": 12
|
||||
}
|
||||
],
|
||||
"error": null,
|
||||
"live_sources": [
|
||||
"K1 - In",
|
||||
"K1 - Out",
|
||||
"K2 - In",
|
||||
"K2 - Out"
|
||||
],
|
||||
"history_rows": 18,
|
||||
"hosts_ok": 2,
|
||||
"hosts_total": 2,
|
||||
"groups_total": 2,
|
||||
"source_key": "kpc"
|
||||
},
|
||||
"updated_at": "2026-09-02T01:00:03Z"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{}
|
||||
@@ -1,2 +1,116 @@
|
||||
# karung-web-admin
|
||||
|
||||
Admin dashboard that aggregates **KTC** (tuang) and **KPC** (masuk) counter data and exposes a single combined API.
|
||||
|
||||
## Sources
|
||||
|
||||
| Role | Project | Typical port | Cameras |
|
||||
|------|---------|--------------|---------|
|
||||
| Tuang | `zenai-ktc-python` | `9000` | `K1 - Use`, `K2 - Use`, ... |
|
||||
| Masuk | `zenai-kpc-python` | `7000` | `K1 - In`, `K1 - Out`, `K2 - In`, `K2 - Out`, ... |
|
||||
|
||||
Host IPs are **not** hardcoded. Configure them per machine in `sources.env`.
|
||||
|
||||
Source URL grouping:
|
||||
|
||||
- `|` separates cameras (`K1`, `K2`, ...)
|
||||
- `,` separates hosts inside one camera (their counts are summed into that `Kn`)
|
||||
- No `|` keeps legacy behavior: each comma-separated host is its own `Kn`
|
||||
- One multi-host camera only: use a trailing `|` (e.g. `http://a,http://b,http://c|`)
|
||||
|
||||
Tuang uses each host's `total_count` as `K{n} - Use`; masuk keeps `K{n} - In` and `K{n} - Out`.
|
||||
|
||||
## 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
|
||||
# K1 = three hosts summed; K2 = one host
|
||||
KTC_BASE_URL=http://10.0.0.11:9000,http://10.0.0.12:9000,http://10.0.0.13:9000|http://10.0.0.14:9000
|
||||
# Same grouping for masuk (In/Out summed per Kn)
|
||||
KPC_BASE_URL=http://10.0.0.11:7000,http://10.0.0.12:7000,http://10.0.0.13:7000|http://10.0.0.14:7000
|
||||
SOURCE_HISTORY_DAYS=30
|
||||
SOURCE_PULL_INTERVAL_SECONDS=30
|
||||
```
|
||||
|
||||
Legacy (unchanged): `http://a:9000,http://b:9000` still means `K1` and `K2` (same for KPC).
|
||||
|
||||
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
|
||||
sudo cp karung-web-admin-pull-sources.service /etc/systemd/system/
|
||||
sudo systemctl daemon-reload
|
||||
sudo systemctl enable --now karung-web-admin-pull-sources.service
|
||||
```
|
||||
|
||||
It loads `/opt/karung-web-admin/sources.env`, runs `pull_sources.py --loop`, and refreshes every `SOURCE_PULL_INTERVAL_SECONDS` (default 30).
|
||||
|
||||
Writes:
|
||||
|
||||
- Live JSON: `karung_tuang.json`, `karung_masuk.json`
|
||||
- History SQLite: `karung_tuang.db`, `karung_masuk.db`
|
||||
- Status: `.pull_status.json`
|
||||
|
||||
One-shot (optional / debugging):
|
||||
|
||||
```bash
|
||||
python pull_sources.py
|
||||
```
|
||||
|
||||
## Web app
|
||||
|
||||
```bash
|
||||
python app.py
|
||||
```
|
||||
|
||||
Or enable the dashboard service:
|
||||
|
||||
```bash
|
||||
sudo cp karung-web-admin.service /etc/systemd/system/
|
||||
sudo systemctl daemon-reload
|
||||
sudo systemctl enable --now karung-web-admin.service
|
||||
```
|
||||
|
||||
If you previously used `frigate-counter_karung-web-admin.service`, disable/remove it:
|
||||
|
||||
```bash
|
||||
sudo systemctl disable --now frigate-counter_karung-web-admin.service
|
||||
sudo rm /etc/systemd/system/frigate-counter_karung-web-admin.service
|
||||
sudo systemctl daemon-reload
|
||||
```
|
||||
|
||||
Dashboard: `http://0.0.0.0:8090/`
|
||||
Combined API: `GET /api/combined?history_days=30`
|
||||
|
||||
## Environment
|
||||
|
||||
Configured in `sources.env` (see `sources.env.example`):
|
||||
|
||||
| Variable | Description |
|
||||
|----------|-------------|
|
||||
| `KTC_BASE_URL` | Tuang/KTC source groups (`\|` = cameras, `,` = hosts in a camera) |
|
||||
| `KPC_BASE_URL` | Masuk/KPC source groups (same grouping syntax) |
|
||||
| `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`) |
|
||||
@@ -1 +0,0 @@
|
||||
COUNTER_UUID=9f556d2a-af5c-41b0-bcd0-6cea0dfa0fed
|
||||
+48
-7
@@ -14,14 +14,55 @@ DEFAULT_LOCAL_TIMEZONE = "Asia/Jakarta"
|
||||
DEFAULT_TUANG_CUTOFF_TIME = "17:00:00"
|
||||
DEFAULT_MASUK_CUTOFF_TIME = "24:00:00"
|
||||
COUNTER_SETTINGS_FILE_NAME = "counter_settings.json"
|
||||
ENV_LOCAL_TIMEZONE = "LOCAL_TIMEZONE"
|
||||
ENV_TUANG_CUTOFF_TIME = "TUANG_CUTOFF_TIME"
|
||||
ENV_MASUK_CUTOFF_TIME = "MASUK_CUTOFF_TIME"
|
||||
|
||||
|
||||
def get_env_counter_settings():
|
||||
"""Return counter settings explicitly set in the environment."""
|
||||
settings = {}
|
||||
|
||||
local_timezone = os.environ.get(ENV_LOCAL_TIMEZONE)
|
||||
if local_timezone:
|
||||
settings["local_timezone"] = normalize_timezone(local_timezone)
|
||||
|
||||
tuang_cutoff_time = os.environ.get(ENV_TUANG_CUTOFF_TIME)
|
||||
if tuang_cutoff_time:
|
||||
settings["tuang_cutoff_time"] = normalize_cutoff_time(tuang_cutoff_time)
|
||||
|
||||
masuk_cutoff_time = os.environ.get(ENV_MASUK_CUTOFF_TIME)
|
||||
if masuk_cutoff_time:
|
||||
settings["masuk_cutoff_time"] = normalize_cutoff_time(masuk_cutoff_time)
|
||||
|
||||
return settings
|
||||
|
||||
|
||||
def apply_env_counter_settings(settings):
|
||||
merged_settings = dict(settings or {})
|
||||
merged_settings.update(get_env_counter_settings())
|
||||
return merged_settings
|
||||
|
||||
|
||||
def get_env_locked_counter_settings():
|
||||
locked_settings = []
|
||||
if os.environ.get(ENV_LOCAL_TIMEZONE):
|
||||
locked_settings.append("local_timezone")
|
||||
if os.environ.get(ENV_TUANG_CUTOFF_TIME):
|
||||
locked_settings.append("tuang_cutoff_time")
|
||||
if os.environ.get(ENV_MASUK_CUTOFF_TIME):
|
||||
locked_settings.append("masuk_cutoff_time")
|
||||
return locked_settings
|
||||
|
||||
|
||||
def default_counter_settings():
|
||||
return {
|
||||
"local_timezone": DEFAULT_LOCAL_TIMEZONE,
|
||||
"tuang_cutoff_time": DEFAULT_TUANG_CUTOFF_TIME,
|
||||
"masuk_cutoff_time": DEFAULT_MASUK_CUTOFF_TIME,
|
||||
}
|
||||
return apply_env_counter_settings(
|
||||
{
|
||||
"local_timezone": DEFAULT_LOCAL_TIMEZONE,
|
||||
"tuang_cutoff_time": DEFAULT_TUANG_CUTOFF_TIME,
|
||||
"masuk_cutoff_time": DEFAULT_MASUK_CUTOFF_TIME,
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def get_settings_file_candidates(base_dir=None):
|
||||
@@ -38,8 +79,8 @@ def get_settings_file_candidates(base_dir=None):
|
||||
current_dir.parent / COUNTER_SETTINGS_FILE_NAME,
|
||||
current_dir.parent / "karung-web-admin" / COUNTER_SETTINGS_FILE_NAME,
|
||||
current_dir.parent / "karung-web" / COUNTER_SETTINGS_FILE_NAME,
|
||||
Path("/etc/frigate-counter/karung-web-admin") / COUNTER_SETTINGS_FILE_NAME,
|
||||
Path("/etc/frigate-counter/karung-web") / COUNTER_SETTINGS_FILE_NAME,
|
||||
Path("/opt/karung-web-admin") / COUNTER_SETTINGS_FILE_NAME,
|
||||
Path("/opt/karung-web") / COUNTER_SETTINGS_FILE_NAME,
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
+217
@@ -0,0 +1,217 @@
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
from datetime import datetime
|
||||
|
||||
DISPLAY_DATE_FORMAT = '%d-%m-%Y'
|
||||
K_CAMERA_PATTERN = re.compile(r'^K(\d+)\b', re.IGNORECASE)
|
||||
KANDANG_NUMBER_PATTERN = re.compile(r'kandang\s*(\d+)\s*$', re.IGNORECASE)
|
||||
|
||||
|
||||
def normalize_base_url(raw_url):
|
||||
url = str(raw_url or '').strip().rstrip('/')
|
||||
if not url:
|
||||
return 'http://localhost:5002'
|
||||
if not url.startswith(('http://', 'https://')):
|
||||
url = f'http://{url}'
|
||||
return url
|
||||
|
||||
|
||||
def get_cycle_api_config():
|
||||
base_url = normalize_base_url(os.environ.get('CYCLE_API_BASE_URL', 'http://localhost:5002'))
|
||||
api_key = os.environ.get('CYCLE_API_KEY', '').strip()
|
||||
return base_url, api_key
|
||||
|
||||
|
||||
def is_cycle_api_configured():
|
||||
_, api_key = get_cycle_api_config()
|
||||
return bool(api_key)
|
||||
|
||||
|
||||
def fetch_paginated_results(path, *, timeout=15):
|
||||
base_url, api_key = get_cycle_api_config()
|
||||
if not api_key:
|
||||
raise RuntimeError('CYCLE_API_KEY is not configured.')
|
||||
|
||||
results = []
|
||||
next_url = urllib.parse.urljoin(f'{base_url}/', path.lstrip('/'))
|
||||
|
||||
while next_url:
|
||||
request = urllib.request.Request(next_url, method='GET')
|
||||
request.add_header('Accept', 'application/json')
|
||||
request.add_header('X-API-Key', api_key)
|
||||
|
||||
try:
|
||||
with urllib.request.urlopen(request, timeout=timeout) as response:
|
||||
payload = json.loads(response.read().decode('utf-8'))
|
||||
except urllib.error.HTTPError as error:
|
||||
body = error.read().decode('utf-8', errors='replace')
|
||||
raise RuntimeError(f'Cycle API request failed ({error.code}): {body}') from error
|
||||
except (urllib.error.URLError, TimeoutError, json.JSONDecodeError) as error:
|
||||
raise RuntimeError(f'Cycle API request failed: {error}') from error
|
||||
|
||||
if isinstance(payload, dict) and 'results' in payload:
|
||||
page_results = payload.get('results') or []
|
||||
if not isinstance(page_results, list):
|
||||
raise RuntimeError('Cycle API paginated response is invalid.')
|
||||
results.extend(page_results)
|
||||
next_url = payload.get('next')
|
||||
continue
|
||||
|
||||
if isinstance(payload, list):
|
||||
results.extend(payload)
|
||||
next_url = None
|
||||
continue
|
||||
|
||||
raise RuntimeError('Cycle API response format is not supported.')
|
||||
|
||||
return results
|
||||
|
||||
|
||||
def extract_k_index(source_name):
|
||||
match = K_CAMERA_PATTERN.match(str(source_name or '').strip())
|
||||
if not match:
|
||||
return None
|
||||
|
||||
return int(match.group(1))
|
||||
|
||||
|
||||
def extract_kandang_number(kandang_name):
|
||||
match = KANDANG_NUMBER_PATTERN.search(str(kandang_name or '').strip())
|
||||
if not match:
|
||||
return None
|
||||
|
||||
return int(match.group(1))
|
||||
|
||||
|
||||
def find_kandang_by_k_index(kandangs, k_index):
|
||||
if k_index is None:
|
||||
return None
|
||||
|
||||
for kandang in kandangs:
|
||||
kandang_name = str(kandang.get('kandang_name', '')).strip()
|
||||
number_match = KANDANG_NUMBER_PATTERN.search(kandang_name)
|
||||
if number_match and int(number_match.group(1)) == k_index:
|
||||
return kandang
|
||||
|
||||
return None
|
||||
|
||||
|
||||
def format_dashboard_date(iso_date):
|
||||
if not iso_date:
|
||||
return ''
|
||||
|
||||
return datetime.strptime(iso_date, '%Y-%m-%d').strftime(DISPLAY_DATE_FORMAT)
|
||||
|
||||
|
||||
def select_cycle_for_kandang(cycles, kandang_id):
|
||||
kandang_cycles = [
|
||||
cycle for cycle in cycles
|
||||
if cycle.get('kandang') == kandang_id
|
||||
]
|
||||
if not kandang_cycles:
|
||||
return None
|
||||
|
||||
active_cycles = [cycle for cycle in kandang_cycles if cycle.get('status') == 'active']
|
||||
candidate_cycles = active_cycles or kandang_cycles
|
||||
return max(candidate_cycles, key=lambda cycle: cycle.get('id', 0))
|
||||
|
||||
|
||||
def format_cycle_name(cycle_start, cycle_end):
|
||||
start = str(cycle_start or '').strip()
|
||||
end = str(cycle_end or '').strip()
|
||||
if start and end:
|
||||
return f'{start} - {end}'
|
||||
if start:
|
||||
return f'{start} - akhir data'
|
||||
if end:
|
||||
return f'awal data - {end}'
|
||||
return 'Siklus Aktif'
|
||||
|
||||
|
||||
def build_cycle_settings(cycle, kandang):
|
||||
kandang_name = kandang.get('kandang_name', f"Kandang {kandang.get('id')}")
|
||||
cycle_start = format_dashboard_date(cycle.get('start_date'))
|
||||
cycle_end = format_dashboard_date(cycle.get('end_date'))
|
||||
return {
|
||||
'id': cycle.get('id', ''),
|
||||
'kandang_id': kandang.get('id'),
|
||||
'k_index': extract_kandang_number(kandang_name),
|
||||
'kandang_name': kandang_name,
|
||||
'cycle_name': format_cycle_name(cycle_start, cycle_end),
|
||||
'cycle_start': cycle_start,
|
||||
'cycle_end': cycle_end,
|
||||
'saldo_awal': cycle.get('feed_initial_balance') or 0,
|
||||
'feed_initial_balance_date': format_dashboard_date(cycle.get('feed_initial_balance_date')),
|
||||
'status': cycle.get('status', ''),
|
||||
'flock': cycle.get('flock'),
|
||||
}
|
||||
|
||||
|
||||
def build_all_cycles_by_kandang(kandangs, cycles):
|
||||
cycles_by_kandang = {}
|
||||
|
||||
for kandang in kandangs:
|
||||
kandang_id = kandang.get('id')
|
||||
if kandang_id is None:
|
||||
continue
|
||||
|
||||
kandang_cycles = [
|
||||
cycle for cycle in cycles
|
||||
if cycle.get('kandang') == kandang_id
|
||||
]
|
||||
kandang_cycles.sort(key=lambda cycle: cycle.get('id', 0), reverse=True)
|
||||
cycles_by_kandang[kandang_id] = [
|
||||
build_cycle_settings(cycle, kandang)
|
||||
for cycle in kandang_cycles
|
||||
]
|
||||
|
||||
return cycles_by_kandang
|
||||
|
||||
|
||||
def load_kandang_cycle_context():
|
||||
kandangs = fetch_paginated_results('/api/v1/kandangs/')
|
||||
cycles = fetch_paginated_results('/api/v1/cycles/')
|
||||
|
||||
kandang_by_k_index = {}
|
||||
cycle_by_kandang_id = {}
|
||||
kandang_cycles = []
|
||||
|
||||
for kandang in kandangs:
|
||||
k_index = extract_kandang_number(kandang.get('kandang_name'))
|
||||
if k_index is not None:
|
||||
kandang_by_k_index[k_index] = kandang
|
||||
|
||||
cycle = select_cycle_for_kandang(cycles, kandang.get('id'))
|
||||
if not cycle:
|
||||
continue
|
||||
|
||||
cycle_settings = build_cycle_settings(cycle, kandang)
|
||||
cycle_by_kandang_id[kandang['id']] = cycle_settings
|
||||
kandang_cycles.append(cycle_settings)
|
||||
|
||||
kandang_cycles.sort(
|
||||
key=lambda item: (item.get('k_index') is None, item.get('k_index') or 0, item.get('kandang_id') or 0)
|
||||
)
|
||||
return {
|
||||
'kandangs': kandangs,
|
||||
'kandang_by_k_index': kandang_by_k_index,
|
||||
'cycle_by_kandang_id': cycle_by_kandang_id,
|
||||
'kandang_cycles': kandang_cycles,
|
||||
'all_cycles_by_kandang': build_all_cycles_by_kandang(kandangs, cycles),
|
||||
}
|
||||
|
||||
|
||||
def get_cycle_settings_for_camera(source_name, cycle_context):
|
||||
k_index = extract_k_index(source_name)
|
||||
if k_index is None:
|
||||
return None
|
||||
|
||||
kandang = find_kandang_by_k_index(cycle_context.get('kandangs', []), k_index)
|
||||
if not kandang:
|
||||
return None
|
||||
|
||||
return cycle_context.get('cycle_by_kandang_id', {}).get(kandang['id'])
|
||||
Binary file not shown.
@@ -1,15 +0,0 @@
|
||||
[Unit]
|
||||
Description=Admin Karung Counter Web
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
WorkingDirectory=/etc/frigate-counter/karung-web-admin
|
||||
ExecStart=/usr/bin/python3 -u /etc/frigate-counter/karung-web-admin/app.py
|
||||
Restart=always
|
||||
RestartSec=5s
|
||||
StandardOutput=journal
|
||||
StandardError=journal
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
@@ -0,0 +1,17 @@
|
||||
[Unit]
|
||||
Description=Pull KTC/KPC counters into karung-web-admin
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
WorkingDirectory=/opt/karung-web-admin
|
||||
# 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
|
||||
StandardOutput=journal
|
||||
StandardError=journal
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
@@ -0,0 +1,15 @@
|
||||
[Unit]
|
||||
Description=Karung Web Admin Dashboard
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
WorkingDirectory=/opt/karung-web-admin
|
||||
ExecStart=/usr/bin/python3 -u /opt/karung-web-admin/app.py
|
||||
Restart=always
|
||||
RestartSec=5s
|
||||
StandardOutput=journal
|
||||
StandardError=journal
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
Binary file not shown.
+9
-3
@@ -1,8 +1,14 @@
|
||||
{
|
||||
"kandang_1_karung_masuk": {
|
||||
"K1 - In": {
|
||||
"karung": 160
|
||||
},
|
||||
"K1 - Out": {
|
||||
"karung": 0
|
||||
},
|
||||
"kandang_1_karung_masuk_out": {
|
||||
"K2 - In": {
|
||||
"karung": 80
|
||||
},
|
||||
"K2 - Out": {
|
||||
"karung": 0
|
||||
}
|
||||
}
|
||||
}
|
||||
Binary file not shown.
+3
-3
@@ -1,8 +1,8 @@
|
||||
{
|
||||
"karung_tuang_feeder_kiri_atas": {
|
||||
"K1 - Use": {
|
||||
"karung": 0
|
||||
},
|
||||
"karung_tuang_feeder_kanan_bawah": {
|
||||
"K2 - Use": {
|
||||
"karung": 0
|
||||
}
|
||||
}
|
||||
}
|
||||
+536
@@ -0,0 +1,536 @@
|
||||
#!/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())
|
||||
@@ -0,0 +1,25 @@
|
||||
# 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
|
||||
#
|
||||
# Grouping:
|
||||
# `|` separates cameras (K1, K2, ...)
|
||||
# `,` separates hosts inside one camera (summed into that Kn)
|
||||
# No `|` keeps legacy behavior: each host is its own Kn
|
||||
# One camera from many hosts: add a trailing `|`
|
||||
#
|
||||
# Restart after changes:
|
||||
# sudo systemctl restart karung-web-admin-pull-sources.service
|
||||
|
||||
# Tuang (KTC)
|
||||
# Example: K1 = three hosts summed, K2 = one host
|
||||
# KTC_BASE_URL=http://REPLACE_HOST_A:9000,http://REPLACE_HOST_B:9000,http://REPLACE_HOST_C:9000|http://REPLACE_HOST_D:9000
|
||||
KTC_BASE_URL=http://REPLACE_HOST_A:9000|http://REPLACE_HOST_B:9000
|
||||
|
||||
# Masuk (KPC) — same grouping; hosts in a group sum into Kn - In / Kn - Out
|
||||
# Example: K1 = three hosts summed, K2 = two hosts summed
|
||||
# KPC_BASE_URL=http://REPLACE_HOST_A:7000,http://REPLACE_HOST_B:7000,http://REPLACE_HOST_C:7000|http://REPLACE_HOST_D:7000,http://REPLACE_HOST_E:7000
|
||||
KPC_BASE_URL=http://REPLACE_HOST_A:7000|http://REPLACE_HOST_B:7000
|
||||
|
||||
SOURCE_HISTORY_DAYS=30
|
||||
SOURCE_PULL_INTERVAL_SECONDS=30
|
||||
-284
@@ -1,284 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Sync counter data from SQLite databases to API.
|
||||
|
||||
karung_tuang:
|
||||
- 17:00-23:59: Syncs CURRENT date
|
||||
- 00:00-16:59: Syncs PREVIOUS date (only if not already synced)
|
||||
|
||||
karung_masuk: Syncs PREVIOUS date from 00:00 onwards, only if data differs
|
||||
(skips if already successfully synced for that date).
|
||||
"""
|
||||
|
||||
import os
|
||||
import sqlite3
|
||||
import json
|
||||
import urllib.request
|
||||
import urllib.error
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from counter_time import get_business_date, load_counter_settings
|
||||
|
||||
|
||||
def get_db_path(db_name):
|
||||
"""Get database path in same directory as script."""
|
||||
script_dir = os.path.dirname(os.path.abspath(__file__))
|
||||
return os.path.join(script_dir, db_name)
|
||||
|
||||
|
||||
def get_previous_business_date(counter_type, settings):
|
||||
"""Get the previously completed business date in YYYY-MM-DD format."""
|
||||
return (get_business_date(counter_type, settings) - timedelta(days=1)).strftime("%Y-%m-%d")
|
||||
|
||||
|
||||
def get_sync_state_file():
|
||||
"""Get path to sync state file."""
|
||||
script_dir = os.path.dirname(os.path.abspath(__file__))
|
||||
return os.path.join(script_dir, ".sync_state.json")
|
||||
|
||||
|
||||
def load_sync_state():
|
||||
"""Load sync state from file."""
|
||||
state_file = get_sync_state_file()
|
||||
if os.path.exists(state_file):
|
||||
try:
|
||||
with open(state_file, "r") as f:
|
||||
return json.load(f)
|
||||
except (json.JSONDecodeError, IOError):
|
||||
return {}
|
||||
return {}
|
||||
|
||||
|
||||
def save_sync_state(state):
|
||||
"""Save sync state to file."""
|
||||
state_file = get_sync_state_file()
|
||||
try:
|
||||
with open(state_file, "w") as f:
|
||||
json.dump(state, f)
|
||||
except IOError as e:
|
||||
print(f"Warning: Could not save sync state: {e}")
|
||||
|
||||
|
||||
def is_already_synced(state, db_name, date):
|
||||
"""Check if database was already synced for given date."""
|
||||
return state.get(db_name, {}).get("last_synced_date") == date
|
||||
|
||||
|
||||
def mark_synced(state, db_name, date):
|
||||
"""Mark database as synced for given date."""
|
||||
if db_name not in state:
|
||||
state[db_name] = {}
|
||||
state[db_name]["last_synced_date"] = date
|
||||
|
||||
|
||||
def get_current_business_date(counter_type, settings):
|
||||
"""Get the active business date in YYYY-MM-DD format."""
|
||||
return get_business_date(counter_type, settings).strftime("%Y-%m-%d")
|
||||
|
||||
|
||||
def get_local_data(db_path, table_name, date):
|
||||
"""Get counter data from SQLite database for a specific date."""
|
||||
conn = sqlite3.connect(db_path)
|
||||
cursor = conn.cursor()
|
||||
|
||||
if table_name == "karung_counts":
|
||||
cursor.execute(
|
||||
"SELECT camera_name, date, counter_value FROM karung_counts WHERE date = ?",
|
||||
(date,),
|
||||
)
|
||||
else:
|
||||
cursor.execute(
|
||||
"SELECT camera_name, date(date) as d, counter_value FROM counter_data WHERE date(date) = ?",
|
||||
(date,),
|
||||
)
|
||||
|
||||
results = cursor.fetchall()
|
||||
conn.close()
|
||||
|
||||
data = {}
|
||||
for camera_name, date_val, counter_value in results:
|
||||
key = f"{camera_name}_{date_val}"
|
||||
data[key] = {
|
||||
"camera_name": camera_name,
|
||||
"date": date_val,
|
||||
"value": counter_value,
|
||||
}
|
||||
return data
|
||||
|
||||
|
||||
def fetch_api_data(uuid):
|
||||
"""Fetch counter data from API."""
|
||||
url = f"https://dashboard.cpsp.id/api/cpsp/counter/{uuid}/"
|
||||
|
||||
req = urllib.request.Request(url, method="GET")
|
||||
req.add_header("Accept", "application/json")
|
||||
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=30) as response:
|
||||
result = json.loads(response.read().decode("utf-8"))
|
||||
if result.get("success") and "data" in result:
|
||||
return result["data"].get("camera_counter", [])
|
||||
except urllib.error.HTTPError as e:
|
||||
print(f"HTTP Error {e.code}: {e.reason}")
|
||||
except urllib.error.URLError as e:
|
||||
print(f"URL Error: {e.reason}")
|
||||
except json.JSONDecodeError as e:
|
||||
print(f"JSON Decode Error: {e}")
|
||||
|
||||
return []
|
||||
|
||||
|
||||
def update_api(uuid, camera_name, date, value):
|
||||
"""Update counter data via API (dummy endpoint)."""
|
||||
# Dummy endpoint - replace with actual endpoint when available
|
||||
url = f"https://dashboard.cpsp.id/api/cpsp/counter/{uuid}/update/"
|
||||
|
||||
payload = {"camera_name": camera_name, "date": date, "value": value}
|
||||
|
||||
data = json.dumps(payload).encode("utf-8")
|
||||
req = urllib.request.Request(url, data=data, method="POST")
|
||||
req.add_header("Content-Type", "application/json")
|
||||
req.add_header("Accept", "application/json")
|
||||
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=30) as response:
|
||||
result = json.loads(response.read().decode("utf-8"))
|
||||
print(f"Updated {camera_name} for {date}: {value}")
|
||||
return result
|
||||
except urllib.error.HTTPError as e:
|
||||
print(f"HTTP Error {e.code} updating {camera_name}: {e.reason}")
|
||||
# For dummy endpoint, just log the intended update
|
||||
print(f"[DUMMY] Would update {camera_name} for {date} with value {value}")
|
||||
return None
|
||||
except urllib.error.URLError as e:
|
||||
print(f"URL Error: {e.reason}")
|
||||
print(f"[DUMMY] Would update {camera_name} for {date} with value {value}")
|
||||
return None
|
||||
|
||||
|
||||
def sync_database(db_config, date, uuid, state=None, skip_if_synced=False):
|
||||
"""Sync a single database with API.
|
||||
|
||||
Args:
|
||||
db_config: Database configuration dict
|
||||
date: Date string to sync
|
||||
uuid: API UUID
|
||||
state: Optional sync state dict
|
||||
skip_if_synced: If True, skip sync if already marked as synced for this date
|
||||
"""
|
||||
db_name = db_config.get("name", "unknown")
|
||||
|
||||
# Check if already synced (for karung_masuk to avoid redundant syncs)
|
||||
if skip_if_synced and state and is_already_synced(state, db_name, date):
|
||||
print(f"Skipping {db_name} - already synced for {date}")
|
||||
return True
|
||||
|
||||
db_path = get_db_path(db_config["db_file"])
|
||||
|
||||
if not os.path.exists(db_path):
|
||||
print(f"Database not found: {db_path}")
|
||||
return False
|
||||
|
||||
# Get local data
|
||||
local_data = get_local_data(db_path, db_config["table"], date)
|
||||
|
||||
if not local_data:
|
||||
print(f"No local data found for {date} in {db_config['db_file']}")
|
||||
return False
|
||||
|
||||
# Get API data
|
||||
api_data = fetch_api_data(uuid)
|
||||
|
||||
# Build API data map
|
||||
api_map = {}
|
||||
for item in api_data:
|
||||
key = f"{item['camera_name']}_{item['date']}"
|
||||
api_map[key] = item["value"]
|
||||
|
||||
has_mismatch = False
|
||||
|
||||
# Compare and sync
|
||||
for key, local_record in local_data.items():
|
||||
camera_name = local_record["camera_name"]
|
||||
local_value = local_record["value"]
|
||||
|
||||
if key not in api_map:
|
||||
print(f"New record: {camera_name} for {date} = {local_value}")
|
||||
update_api(uuid, camera_name, date, local_value)
|
||||
has_mismatch = True
|
||||
elif api_map[key] != local_value:
|
||||
print(
|
||||
f"Mismatch for {camera_name} on {date}: API={api_map[key]}, Local={local_value}"
|
||||
)
|
||||
update_api(uuid, camera_name, date, local_value)
|
||||
has_mismatch = True
|
||||
else:
|
||||
print(f"Synced: {camera_name} for {date} = {local_value}")
|
||||
|
||||
# Mark as synced if no mismatches found (all data matches)
|
||||
if not has_mismatch and state:
|
||||
mark_synced(state, db_name, date)
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def main():
|
||||
"""Main sync function.
|
||||
|
||||
karung_tuang:
|
||||
- 17:00-23:59: Sync CURRENT date
|
||||
- 00:00-16:59: Sync PREVIOUS date (only if not already synced)
|
||||
|
||||
karung_masuk: Sync PREVIOUS date from 00:00 onwards, only if data differs.
|
||||
"""
|
||||
uuid = os.environ.get("COUNTER_UUID")
|
||||
if not uuid:
|
||||
print("Error: COUNTER_UUID environment variable not set")
|
||||
return
|
||||
|
||||
# Load sync state to track which dates have been synced
|
||||
sync_state = load_sync_state()
|
||||
counter_settings = load_counter_settings()
|
||||
|
||||
# Database configurations
|
||||
db_configs = {
|
||||
"karung_tuang": {
|
||||
"name": "karung_tuang",
|
||||
"counter_type": "tuang",
|
||||
"db_file": "karung_tuang.db",
|
||||
"table": "counter_data",
|
||||
"skip_if_synced": False,
|
||||
},
|
||||
"karung_masuk": {
|
||||
"name": "karung_masuk",
|
||||
"counter_type": "masuk",
|
||||
"db_file": "karung_masuk.db",
|
||||
"table": "karung_counts",
|
||||
"skip_if_synced": True, # Skip if already synced for this date
|
||||
},
|
||||
}
|
||||
|
||||
print(f"Timezone: {counter_settings['local_timezone']}")
|
||||
|
||||
# Sync the business date that has most recently closed for each counter.
|
||||
for name, config in db_configs.items():
|
||||
sync_date = get_previous_business_date(config["counter_type"], counter_settings)
|
||||
active_date = get_current_business_date(config["counter_type"], counter_settings)
|
||||
|
||||
print(f"\nSyncing {name} for completed business date {sync_date} (active date {active_date})...")
|
||||
sync_database(
|
||||
config,
|
||||
sync_date,
|
||||
uuid,
|
||||
state=sync_state,
|
||||
skip_if_synced=config.get("skip_if_synced", False),
|
||||
)
|
||||
|
||||
# Save sync state
|
||||
save_sync_state(sync_state)
|
||||
|
||||
print("\nSync complete.")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
+486
-394
File diff suppressed because it is too large.
Load diff
Reference in new issue
Block a user