209 lines
5.8 KiB
Python
209 lines
5.8 KiB
Python
"""SNEK Dashboard — Flask backend receiving SNEK callbacks."""
|
|
|
|
import json
|
|
import os
|
|
import queue
|
|
import sqlite3
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
|
|
from flask import Flask, jsonify, request, send_from_directory, stream_with_context, Response
|
|
|
|
from sensit_decoder import decode as sensit_decode
|
|
|
|
DB_PATH = os.environ.get("DB_PATH", "/data/messages.db")
|
|
Path(DB_PATH).parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
TOKEN = os.environ.get("DASHBOARD_TOKEN", "").strip()
|
|
# Public: HTML shell + static assets (login page needs to load) + webhook + health
|
|
PUBLIC_ROUTES = {"/", "/webhook", "/healthz"}
|
|
PUBLIC_PREFIXES = ("/static/",)
|
|
|
|
app = Flask(__name__, static_folder="static", static_url_path="/static")
|
|
|
|
|
|
@app.before_request
|
|
def _check_token():
|
|
p = request.path
|
|
if not TOKEN or p in PUBLIC_ROUTES or p.startswith(PUBLIC_PREFIXES):
|
|
return None
|
|
supplied = (
|
|
request.headers.get("Authorization", "").removeprefix("Bearer ").strip()
|
|
or request.args.get("token", "").strip()
|
|
)
|
|
if supplied != TOKEN:
|
|
return jsonify({"error": "unauthorized"}), 401
|
|
return None
|
|
|
|
# --- SSE broadcast ---
|
|
_subscribers: list[queue.Queue] = []
|
|
_subscribers_lock = threading.Lock()
|
|
|
|
|
|
def _broadcast(event: dict) -> None:
|
|
payload = json.dumps(event)
|
|
with _subscribers_lock:
|
|
dead = []
|
|
for q in _subscribers:
|
|
try:
|
|
q.put_nowait(payload)
|
|
except queue.Full:
|
|
dead.append(q)
|
|
for q in dead:
|
|
_subscribers.remove(q)
|
|
|
|
|
|
# --- Database ---
|
|
def _init_db() -> None:
|
|
with sqlite3.connect(DB_PATH) as conn:
|
|
conn.execute(
|
|
"""CREATE TABLE IF NOT EXISTS messages (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
ts REAL NOT NULL,
|
|
device TEXT NOT NULL,
|
|
raw TEXT NOT NULL,
|
|
decoded TEXT NOT NULL,
|
|
rssi REAL,
|
|
snr REAL,
|
|
seq_number INTEGER
|
|
)"""
|
|
)
|
|
conn.execute("CREATE INDEX IF NOT EXISTS idx_device_ts ON messages(device, ts DESC)")
|
|
conn.commit()
|
|
|
|
|
|
_init_db()
|
|
|
|
|
|
# --- Routes ---
|
|
@app.route("/")
|
|
def index():
|
|
return send_from_directory("static", "index.html")
|
|
|
|
|
|
def _num(x, cast=float):
|
|
try:
|
|
return cast(x)
|
|
except (TypeError, ValueError):
|
|
return None
|
|
|
|
|
|
@app.route("/webhook", methods=["POST", "GET"])
|
|
def webhook():
|
|
"""SNEK callback endpoint. Accepts JSON body OR query string / form params."""
|
|
src = request.get_json(force=True, silent=True) or request.values
|
|
|
|
device = (src.get("device") or "").lower()
|
|
raw = (src.get("data") or "").lower()
|
|
ts = _num(src.get("time")) or time.time()
|
|
rssi = _num(src.get("rssi"))
|
|
snr = _num(src.get("snr"))
|
|
seq_number = _num(src.get("seqNumber") or src.get("seqnumber"), int)
|
|
|
|
if not device or not raw:
|
|
return jsonify({"error": "missing device or data"}), 400
|
|
|
|
decoded = sensit_decode(raw)
|
|
|
|
with sqlite3.connect(DB_PATH) as conn:
|
|
conn.execute(
|
|
"INSERT INTO messages (ts, device, raw, decoded, rssi, snr, seq_number) "
|
|
"VALUES (?,?,?,?,?,?,?)",
|
|
(ts, device, raw, json.dumps(decoded), rssi, snr, seq_number),
|
|
)
|
|
conn.commit()
|
|
|
|
_broadcast({
|
|
"type": "message",
|
|
"ts": ts,
|
|
"device": device,
|
|
"raw": raw,
|
|
"decoded": decoded,
|
|
"rssi": rssi,
|
|
"snr": snr,
|
|
"seqNumber": seq_number,
|
|
})
|
|
return jsonify({"status": "ok"})
|
|
|
|
|
|
@app.route("/api/devices")
|
|
def api_devices():
|
|
with sqlite3.connect(DB_PATH) as conn:
|
|
cur = conn.execute(
|
|
"SELECT device, COUNT(*) c, MAX(ts) last_ts FROM messages GROUP BY device ORDER BY last_ts DESC"
|
|
)
|
|
rows = cur.fetchall()
|
|
return jsonify([{"device": r[0], "count": r[1], "last_ts": r[2]} for r in rows])
|
|
|
|
|
|
@app.route("/api/messages")
|
|
def api_messages():
|
|
device = (request.args.get("device") or "").lower()
|
|
limit = int(request.args.get("limit") or 100)
|
|
with sqlite3.connect(DB_PATH) as conn:
|
|
if device:
|
|
cur = conn.execute(
|
|
"SELECT ts, device, raw, decoded, rssi, snr, seq_number FROM messages "
|
|
"WHERE device = ? ORDER BY ts DESC LIMIT ?",
|
|
(device, limit),
|
|
)
|
|
else:
|
|
cur = conn.execute(
|
|
"SELECT ts, device, raw, decoded, rssi, snr, seq_number FROM messages "
|
|
"ORDER BY ts DESC LIMIT ?",
|
|
(limit,),
|
|
)
|
|
rows = cur.fetchall()
|
|
return jsonify([
|
|
{
|
|
"ts": r[0],
|
|
"device": r[1],
|
|
"raw": r[2],
|
|
"decoded": json.loads(r[3]),
|
|
"rssi": r[4],
|
|
"snr": r[5],
|
|
"seqNumber": r[6],
|
|
}
|
|
for r in rows
|
|
])
|
|
|
|
|
|
@app.route("/api/decode")
|
|
def api_decode():
|
|
"""Manual decode utility: /api/decode?hex=b60dc86e"""
|
|
h = request.args.get("hex", "")
|
|
return jsonify(sensit_decode(h))
|
|
|
|
|
|
@app.route("/events")
|
|
def events():
|
|
"""SSE stream of live messages."""
|
|
def stream():
|
|
q: queue.Queue = queue.Queue(maxsize=100)
|
|
with _subscribers_lock:
|
|
_subscribers.append(q)
|
|
try:
|
|
yield "event: ping\ndata: connected\n\n"
|
|
while True:
|
|
try:
|
|
payload = q.get(timeout=25)
|
|
yield f"data: {payload}\n\n"
|
|
except queue.Empty:
|
|
yield ": keepalive\n\n"
|
|
finally:
|
|
with _subscribers_lock:
|
|
if q in _subscribers:
|
|
_subscribers.remove(q)
|
|
|
|
return Response(stream_with_context(stream()), mimetype="text/event-stream")
|
|
|
|
|
|
@app.route("/healthz")
|
|
def healthz():
|
|
return jsonify({"status": "ok"})
|
|
|
|
|
|
if __name__ == "__main__":
|
|
app.run(host="0.0.0.0", port=3000)
|