Files
bytetrack-counter-dashboard/recounting_dashboard_upload.py
T
proitlab fccc285ab3 Add OUTPUT_DIR browser + copy-on-load with progress to upload dashboard
LOAD MP4 on an OUTPUT_DIR file copies it into UPLOAD_DIR (chunked,
progress via /api/copy-progress) before spawning RECOUNT_CMD; collision
in UPLOAD_DIR fails with 400. New read-only folder-grouped browser
mirrors recounting_dashboard.py. Docs updated.
2026-08-13 13:33:26 +07:00

540 lines
17 KiB
Python

#!/usr/bin/env python3
"""
Recounting upload dashboard — upload MP4 then launch recount via commander.
Consumes live counter + recount APIs, proxies recount live-video for preview.
"""
import json
import logging
import os
import re
import shlex
import shutil
import signal
import subprocess
import threading
import time
import requests
from datetime import datetime
from pathlib import Path
from flask import Flask, render_template, jsonify, request, Response, send_file
from werkzeug.serving import WSGIRequestHandler
from werkzeug.utils import secure_filename
from dotenv import load_dotenv
load_dotenv()
log = logging.getLogger('werkzeug')
log.setLevel(logging.ERROR)
app = Flask(__name__, template_folder="templates")
app.config["SECRET_KEY"] = os.getenv("SECRET_KEY", "change-me-in-production")
app.config["MAX_CONTENT_LENGTH"] = 512 * 1024 * 1024 # 512 MB
REPO_ROOT = Path(__file__).resolve().parent
UPLOAD_DIR = Path(os.getenv("UPLOAD_DIR", str(REPO_ROOT / "uploads")))
OUTPUT_DIR = Path(os.getenv("OUTPUT_DIR", "/opt/bytetrack-counter"))
LIVE_API_URL = os.getenv("LIVE_API_URL", "http://localhost:5000")
RECOUNT_API_URL = os.getenv("RECOUNT_API_URL", "http://localhost:5001")
SITE_NAME = os.getenv("SITE_NAME", "RECOUNT")
TEMPLATE = os.getenv("RECOUNTING_UPLOAD_TEMPLATE", "recounting_upload_lamborghini.html")
LIVE_DASHBOARD_URL = os.getenv("LIVE_DASHBOARD_URL", "")
DASHBOARD_PORT = int(os.getenv("RECOUNTING_UPLOAD_PORT", "5003"))
DASHBOARD_HOST = os.getenv("DASHBOARD_HOST", "0.0.0.0")
FLASK_DEBUG = os.getenv("FLASK_DEBUG", "false").lower() == "true"
RECOUNT_CMD = os.getenv("RECOUNT_CMD", "bytetrack-counter config.env --source {path}")
SHM_DIR = os.getenv("SHM_DIR", "/dev/shm/bytetrack-counter")
_CONTINUE_FLAG = Path(f"{SHM_DIR}/.continue")
_SWEEP_GRACE_SEC = 5
_SWEEP_MAX_RETRIES = 3
_http_session = requests.Session()
_http_session.timeout = 3
_FILENAME_RE = re.compile(r"batch[_-]?(\d+)[_-](\d{8})[_-]\d{6}", re.IGNORECASE)
_recount_proc = None
_recount_lock = threading.Lock()
_current_file = None
_copy_progress = {}
_copy_lock = threading.Lock()
def _parse_filename(filename):
m = _FILENAME_RE.search(filename)
if m:
raw = m.group(2)
return int(m.group(1)), f"{raw[:4]}-{raw[4:6]}-{raw[6:8]}"
return None, None
def _api_get(base_url, path, default=None):
try:
resp = _http_session.get(f"{base_url}{path}")
if resp.status_code == 200:
return resp.json()
except Exception:
pass
return default
def _stop_recount_locked():
global _recount_proc
proc = _recount_proc
_recount_proc = None
if proc is None:
return
try:
pgid = os.getpgid(proc.pid)
except (ProcessLookupError, PermissionError):
return
try:
os.killpg(pgid, signal.SIGTERM)
except (ProcessLookupError, PermissionError):
pass
try:
proc.wait(timeout=5)
except subprocess.TimeoutExpired:
try:
os.killpg(pgid, signal.SIGKILL)
except (ProcessLookupError, PermissionError):
pass
try:
proc.wait(timeout=3)
except Exception:
pass
except Exception:
pass
def _recount_binary_name():
if not RECOUNT_CMD:
return "bytetrack-counter-cpp"
try:
first = shlex.split(RECOUNT_CMD)[0]
except (ValueError, IndexError):
return "bytetrack-counter-cpp"
return Path(first).name or "bytetrack-counter-cpp"
def _iter_matching_pids(binary):
for entry in os.listdir("/proc"):
if not entry.isdigit():
continue
try:
with open(f"/proc/{entry}/cmdline", "rb") as f:
raw = f.read().split(b"\x00")
except (FileNotFoundError, PermissionError, ProcessLookupError):
continue
args = [t.decode(errors="ignore") for t in raw if t]
for token in args:
for piece in token.split():
if Path(piece).name == binary:
yield int(entry)
break
else:
continue
break
def _sweep_recount_processes(binary):
killed = 0
for _ in range(_SWEEP_MAX_RETRIES):
pids = list(_iter_matching_pids(binary))
if not pids:
break
for pid in pids:
if pid not in list(_iter_matching_pids(binary)):
continue
try:
os.kill(pid, signal.SIGTERM)
killed += 1
except (ProcessLookupError, PermissionError):
pass
time.sleep(_SWEEP_GRACE_SEC)
for pid in list(_iter_matching_pids(binary)):
try:
os.kill(pid, signal.SIGKILL)
except (ProcessLookupError, PermissionError):
pass
time.sleep(1)
return killed, len(list(_iter_matching_pids(binary)))
def _folder_date(filepath, base):
rel = Path(filepath).resolve().relative_to(base.resolve())
parts = rel.parts
if len(parts) >= 2:
return parts[0]
return None
@app.route("/")
def index():
return render_template(TEMPLATE, site_name=SITE_NAME, upload_dir=str(UPLOAD_DIR), output_dir=str(OUTPUT_DIR), live_dashboard_url=LIVE_DASHBOARD_URL)
@app.route("/api/live-progress")
def api_live_progress():
data = _api_get(LIVE_API_URL, "/api/current-batch")
if data and data.get("success"):
return jsonify(data)
return jsonify({"success": False, "count": 0, "batch_number": None, "error": "Live unreachable"}), 200
@app.route("/api/recount-progress")
def api_recount_progress():
data = _api_get(RECOUNT_API_URL, "/api/current-batch")
if data and data.get("success"):
return jsonify(data)
return jsonify({"success": False, "count": 0, "batch_number": None, "error": "Recount unreachable"}), 200
@app.route("/api/uploads")
def api_uploads():
files = []
Path(UPLOAD_DIR).mkdir(parents=True, exist_ok=True)
for f in sorted(Path(UPLOAD_DIR).glob("*.mp4"), key=lambda p: p.stat().st_mtime, reverse=True):
st = f.stat()
files.append({
"name": f.name,
"path": str(f),
"size": st.st_size,
"mtime": datetime.fromtimestamp(st.st_mtime).isoformat(),
})
return jsonify(files)
@app.route("/api/mp4-files")
def api_mp4_files():
folders = {}
output = Path(OUTPUT_DIR)
if not output.exists():
return jsonify([])
for d in sorted(output.iterdir(), reverse=True):
if not d.is_dir():
continue
folder_name = d.name
mp4s = sorted(d.glob("*.mp4"), key=lambda p: _parse_filename(p.name)[0] or 0)
if not mp4s:
continue
files = []
for f in mp4s:
st = f.stat()
files.append({
"name": f.name,
"path": str(f),
"size": st.st_size,
"mtime": datetime.fromtimestamp(st.st_mtime).isoformat(),
})
folders[folder_name] = files
result = []
for folder, files in sorted(folders.items(), reverse=True):
result.append({"folder": folder, "files": files})
return jsonify(result)
@app.route("/api/copy-progress")
def api_copy_progress():
with _copy_lock:
p = dict(_copy_progress)
if not p or not p.get("active"):
return jsonify({"active": False})
total = p.get("total") or 0
done = p.get("done") or 0
pct = round(done * 100.0 / total) if total else 0
return jsonify({
"active": True,
"name": p.get("name"),
"total": total,
"done": done,
"pct": pct,
})
@app.route("/api/upload", methods=["POST"])
def api_upload():
if "file" not in request.files:
return jsonify({"success": False, "error": "No file field"}), 400
f = request.files["file"]
if not f.filename:
return jsonify({"success": False, "error": "No file selected"}), 400
if not f.filename.lower().endswith(".mp4"):
return jsonify({"success": False, "error": "Only .mp4 files allowed"}), 400
Path(UPLOAD_DIR).mkdir(parents=True, exist_ok=True)
dest_name = secure_filename(f.filename)
dest_path = Path(UPLOAD_DIR) / dest_name
f.save(str(dest_path))
st = dest_path.stat()
return jsonify({
"success": True,
"name": dest_name,
"path": str(dest_path),
"size": st.st_size,
})
@app.route("/api/delete-upload", methods=["POST"])
def api_delete_upload():
data = request.get_json(force=True) or {}
fpath = data.get("path", "")
if not fpath or not os.path.isfile(fpath):
return jsonify({"success": False, "error": "File not found"}), 404
f = Path(fpath)
if f.parent.resolve() != Path(UPLOAD_DIR).resolve():
return jsonify({"success": False, "error": "Not in upload dir"}), 403
f.unlink()
return jsonify({"success": True})
def _reset_recount_api():
try:
resp = _http_session.post(f"{RECOUNT_API_URL}/api/reset", timeout=5)
print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] Recount reset: {RECOUNT_API_URL}/api/reset → {resp.status_code}")
except requests.exceptions.Timeout:
print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] Recount reset: {RECOUNT_API_URL}/api/reset → timeout")
except requests.exceptions.ConnectionError as e:
print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] Recount reset: {RECOUNT_API_URL}/api/reset → unreachable: {e}")
except Exception as e:
print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] Recount reset: {RECOUNT_API_URL}/api/reset → error: {e}")
def _spawn_recount_locked(upload_path):
global _recount_proc
cmd_str = RECOUNT_CMD.replace("{path}", shlex.quote(str(upload_path)))
_recount_proc = subprocess.Popen(cmd_str, shell=True, start_new_session=True)
print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] Recount loaded: {cmd_str}")
def _copy_file_with_progress(src, dest):
global _copy_progress
total = os.path.getsize(src)
done = 0
with _copy_lock:
_copy_progress = {
"active": True,
"name": Path(src).name,
"total": total,
"done": 0,
"error": None,
}
try:
with open(src, "rb") as fsrc, open(dest, "wb") as fdst:
while True:
buf = fsrc.read(1024 * 1024)
if not buf:
break
fdst.write(buf)
done += len(buf)
with _copy_lock:
_copy_progress["done"] = done
shutil.copystat(src, dest)
except Exception as e:
with _copy_lock:
_copy_progress["error"] = str(e)
raise
finally:
with _copy_lock:
_copy_progress["active"] = False
@app.route("/api/load-mp4", methods=["POST"])
def load_mp4():
global _current_file
data = request.get_json(force=True) or {}
src_path = data.get("path", "")
if not src_path:
return jsonify({"success": False, "error": "Missing 'path'"}), 400
if not os.path.isfile(src_path):
return jsonify({"success": False, "error": f"File not found: {src_path}"}), 404
src_path = str(Path(src_path).resolve())
upload_path = src_path
if Path(src_path).parent != Path(UPLOAD_DIR).resolve():
dest_path = str(Path(UPLOAD_DIR) / Path(src_path).name)
if os.path.exists(dest_path):
return jsonify({
"success": False,
"error": f"{Path(dest_path).name} already exists in uploads, delete it first",
}), 400
Path(UPLOAD_DIR).mkdir(parents=True, exist_ok=True)
try:
_copy_file_with_progress(src_path, dest_path)
except Exception as e:
with _copy_lock:
err = _copy_progress.get("error") or str(e)
return jsonify({"success": False, "error": f"Copy failed: {err}"}), 500
upload_path = dest_path
with _recount_lock:
_stop_recount_locked()
_current_file = upload_path
try:
_spawn_recount_locked(upload_path)
except Exception as e:
_current_file = None
return jsonify({"success": False, "error": f"Failed to start: {e}"}), 500
_reset_recount_api()
return jsonify({
"success": True,
"stream_url": "/api/proxy-stream",
"file": Path(upload_path).name,
"path": upload_path,
})
@app.route("/api/start-recount", methods=["POST"])
def start_recount():
with _recount_lock:
current = _current_file
if not current:
return jsonify({"success": False, "error": "Load an MP4 first"}), 400
try:
Path(SHM_DIR).mkdir(parents=True, exist_ok=True)
_CONTINUE_FLAG.touch(exist_ok=True)
print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] Recount continue: {_CONTINUE_FLAG}")
except Exception as e:
return jsonify({"success": False, "error": f"Failed to signal: {e}"}), 500
return jsonify({
"success": True,
"stream_url": "/api/proxy-stream",
"file": Path(current).name,
"path": current,
})
@app.route("/api/stop-recount", methods=["POST"])
def stop_recount():
global _current_file
with _recount_lock:
_stop_recount_locked()
_current_file = None
killed, remaining = _sweep_recount_processes(_recount_binary_name())
stamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
if remaining:
print(f"[{stamp}] Recount stop: killed {killed}, remaining {remaining} (SIGKILL could not reap)")
else:
print(f"[{stamp}] Recount stop: killed {killed}, remaining {remaining}")
return jsonify({"success": True, "killed": killed, "remaining": remaining})
@app.route("/api/state")
def api_state():
with _recount_lock:
running = _recount_proc is not None and _recount_proc.poll() is None
current = _current_file
result = {
"streaming": running,
}
if current:
result["file"] = current
result["file_name"] = Path(current).name
result["stream_url"] = "/api/proxy-stream"
batch_num, date_str = _parse_filename(Path(current).name)
if batch_num is not None:
result["batch_number"] = batch_num
result["batch_date"] = date_str
return jsonify(result)
@app.route("/api/proxy-stream")
def api_proxy_stream():
try:
resp = _http_session.get(f"{RECOUNT_API_URL}/api/live-video", stream=True, timeout=5)
if resp.status_code != 200:
return Response("stream unavailable", status=502)
def generate():
for chunk in resp.iter_content(chunk_size=8192):
if chunk:
yield chunk
resp.close()
return Response(
generate(),
mimetype=resp.headers.get("Content-Type", "multipart/x-mixed-replace; boundary=frame"),
)
except Exception:
return Response("stream unavailable", status=502)
@app.route("/api/batch-result")
def api_batch_result():
fpath = request.args.get("file", "")
filename = Path(fpath).name
batch_num, date_str = _parse_filename(filename)
if batch_num is None:
return jsonify({"success": False, "error": "Cannot parse batch number from filename"}), 200
if not date_str:
date_str = _folder_date(fpath, Path(UPLOAD_DIR))
if not date_str:
return jsonify({"success": False, "error": "Cannot determine date from filename"}), 200
def _match_in_batches(batches):
for b in batches:
if b.get("batch_number") == batch_num:
return b
return None
data = _api_get(LIVE_API_URL, f"/api/day-detail/{date_str}")
if data:
b = _match_in_batches(data.get("batches", []))
if b:
return jsonify({
"success": True,
"batch_number": batch_num,
"date": date_str,
"count": b["count"],
"start_time": b.get("start_time"),
"end_time": b.get("end_time"),
})
recent = _api_get(LIVE_API_URL, "/api/recent-batches?limit=100")
if recent:
for entry in recent if isinstance(recent, list) else recent.get("batches", []):
if isinstance(entry, dict) and entry.get("batch_number") == batch_num:
return jsonify({
"success": True,
"batch_number": batch_num,
"date": entry.get("date", ""),
"count": entry["count"],
"start_time": entry.get("start_time"),
"end_time": entry.get("end_time"),
})
return jsonify({"success": False, "error": f"Batch #{batch_num} not found"}), 200
@app.route("/api/download-mp4")
def api_download_mp4():
fpath = request.args.get("file", "")
if not fpath or not os.path.isfile(fpath):
return jsonify({"success": False, "error": "File not found"}), 404
return send_file(fpath, as_attachment=True, download_name=Path(fpath).name)
if __name__ == "__main__":
WSGIRequestHandler.protocol_version = "HTTP/1.1"
Path(UPLOAD_DIR).mkdir(parents=True, exist_ok=True)
print(f"Recounting upload dashboard at http://{DASHBOARD_HOST}:{DASHBOARD_PORT}")
print(f"Live API: {LIVE_API_URL}")
print(f"Recount API: {RECOUNT_API_URL}")
print(f"Upload dir: {UPLOAD_DIR}")
print(f"Output dir: {OUTPUT_DIR}")
print(f"Recount command: {RECOUNT_CMD}")
app.run(host=DASHBOARD_HOST, port=DASHBOARD_PORT, debug=FLASK_DEBUG)