From be2d81b3ab69c6b9fac21d20eb5041d8ef7d86b7 Mon Sep 17 00:00:00 2001 From: Mike Wichers Date: Wed, 26 Aug 2026 19:38:33 -0400 Subject: [PATCH] Lead pipeline: local SQLite store as source of truth Form submissions now land in /data/scout73.db before anything leaves the box. The Google Sheet and the ntfy push become mirrors whose per-record outcome is recorded, so a failed sheet write is replayable instead of surviving only as a push telling you to retype it. The table is a generic record store keyed by 'kind', with workflow columns and an audit trail, so RSVPs and other forms can land in the same place later without a migration. Admin panel talks to /api/admin over HTTP and never opens the DB file. - app/store.py schema, writes, reads, one-time leads.jsonl backfill - app/admin_api.py token-gated API; fails closed when ADMIN_TOKEN is unset - app/app.py three hunks: imports, boot init, join_post - compose ADMIN_TOKEN passed through from the Portainer stack env --- app/admin_api.py | 104 +++++++++++++ app/app.py | 26 ++++ app/store.py | 375 +++++++++++++++++++++++++++++++++++++++++++++ docker-compose.yml | 1 + 4 files changed, 506 insertions(+) create mode 100644 app/admin_api.py create mode 100644 app/store.py diff --git a/app/admin_api.py b/app/admin_api.py new file mode 100644 index 0000000..ed551ef --- /dev/null +++ b/app/admin_api.py @@ -0,0 +1,104 @@ +""" +admin_api.py - read/act API over the internal record store. + +This is the seam the future scout admin panel plugs into. The panel talks HTTP +to these endpoints; it never opens the SQLite file directly. That keeps the +panel deployable anywhere (separate container, separate host) and keeps this +app the only writer to its own database. + +Auth: every route requires the X-Admin-Token header to match ADMIN_TOKEN. +If ADMIN_TOKEN is unset the whole router returns 503 - it FAILS CLOSED. These +endpoints expose parent names, emails and phone numbers for minors' families, +so an unconfigured deployment must not serve them. +""" + +import hmac +import os + +from fastapi import APIRouter, Header, HTTPException, Query +from pydantic import BaseModel + +import store + +ADMIN_TOKEN = os.environ.get("ADMIN_TOKEN", "").strip() + +router = APIRouter(prefix="/api/admin", tags=["admin"]) + + +def _auth(token): + if not ADMIN_TOKEN: + raise HTTPException(503, "admin API disabled: ADMIN_TOKEN is not set") + if not token or not hmac.compare_digest(token, ADMIN_TOKEN): + raise HTTPException(401, "bad or missing X-Admin-Token") + + +class RecordPatch(BaseModel): + status: str | None = None + assigned_to: str | None = None + notes: str | None = None + actor: str | None = None + + +@router.get("/summary") +def get_summary(x_admin_token: str = Header(None)): + _auth(x_admin_token) + return store.summary() + + +@router.get("/records") +def get_records(kind: str = None, status: str = None, since: str = None, + q: str = None, limit: int = Query(100, ge=1, le=500), offset: int = 0, + x_admin_token: str = Header(None)): + _auth(x_admin_token) + return {"records": store.list_records(kind=kind, status=status, since=since, + q=q, limit=limit, offset=offset)} + + +@router.get("/records/{record_id}") +def get_one(record_id: str, x_admin_token: str = Header(None)): + _auth(x_admin_token) + rec = store.get_record(record_id, with_audit=True) + if not rec: + raise HTTPException(404, "no such record") + return rec + + +@router.patch("/records/{record_id}") +def patch_one(record_id: str, patch: RecordPatch, x_admin_token: str = Header(None)): + _auth(x_admin_token) + try: + rec = store.update_record(record_id, actor=patch.actor or "admin", + status=patch.status, assigned_to=patch.assigned_to, + notes=patch.notes) + except ValueError as e: + raise HTTPException(422, str(e)) + if not rec: + raise HTTPException(404, "no such record") + return rec + + +@router.get("/mirrors/failed") +def failed(target: str = "google_sheet", x_admin_token: str = Header(None)): + _auth(x_admin_token) + return {"target": target, "records": store.failed_mirror_records(target)} + + +@router.post("/mirrors/retry") +def retry(target: str = "google_sheet", x_admin_token: str = Header(None)): + """Replay records whose copy to an external target failed. Idempotent-ish: + a record already marked ok is never retried.""" + _auth(x_admin_token) + if target != "google_sheet": + raise HTTPException(422, "only google_sheet retry is implemented") + import app as main_app + done, failed_ids = 0, [] + for rec in store.failed_mirror_records(target, limit=200): + try: + main_app.sheet_append(rec["payload"]) + store.set_mirror(rec["id"], target, True) + store.log(rec["id"], "admin", "mirror_retry_ok", target) + done += 1 + except Exception as e: + store.set_mirror(rec["id"], target, False, e) + failed_ids.append(rec["id"]) + return {"retried_ok": done, "still_failing": failed_ids} diff --git a/app/app.py b/app/app.py index f1e492f..4d8a6b4 100644 --- a/app/app.py +++ b/app/app.py @@ -4,6 +4,11 @@ from fastapi import FastAPI, Form from fastapi.responses import HTMLResponse, RedirectResponse from fastapi.staticfiles import StaticFiles +# Internal record store + admin API. The store is the source of truth for +# anything a family submits; see store.py for why. +import store +import admin_api + # --------------------------------------------------------------------------- # Pack & Troop 73 - greenlanescouts73.org # Calendar lives in ONE file: /data/events.json (seeded from app/events.json on @@ -29,6 +34,10 @@ SHEET_HEADERS = ["Timestamp", "Parent Name", "Email", "Phone", app = FastAPI(title="Pack & Troop 73") app.mount("/static", StaticFiles(directory="/app/static"), name="static") +# Creates the schema if absent and backfills the legacy leads.jsonl once. +store.init(backfill_jsonl=LEADS) +app.include_router(admin_api.router) + if not EVENTS_LIVE.exists() and EVENTS_SEED.exists(): try: shutil.copyfile(EVENTS_SEED, EVENTS_LIVE) except Exception: pass @@ -628,13 +637,30 @@ def join_post(parent_name: str = Form(...), email: str = Form(...), phone: str = f.write(json.dumps(rec) + "\n") except Exception: pass + # Land it internally FIRST. Everything below this point is a mirror of a + # row that already exists locally, and each mirror's outcome is recorded + # against the record so a failure is replayable instead of just shouted. + record_id = None + try: + record_id = store.insert_record( + kind="join_lead", source=rec["source"], payload=rec, + unit=store.unit_of(rec["interested_in"]), + display_name=rec["parent_name"], email=rec["email"], + phone=rec["phone"], created_at=rec["ts"]) + except Exception as e: + print("STORE WRITE FAILED for %s <%s>: %s" % (rec["parent_name"], rec["email"], e), flush=True) sheet_error = None try: sheet_append(rec) except Exception as e: sheet_error = str(e) print("SHEET WRITE FAILED for %s <%s>: %s" % (rec["parent_name"], rec["email"], e), flush=True) + if record_id: + store.set_mirror(record_id, "google_sheet", sheet_error is None, sheet_error) delivered = pipeline_notify(rec, sheet_error) + if record_id: + store.set_mirror(record_id, "ntfy", bool(delivered), + None if delivered else "ntfy push not delivered") if sheet_error and not delivered and NTFY_URL: # last-ditch push on the fallback topic so a lead is never silently lost try: diff --git a/app/store.py b/app/store.py new file mode 100644 index 0000000..bd63992 --- /dev/null +++ b/app/store.py @@ -0,0 +1,375 @@ +""" +store.py - internal record store for greenlanescouts73.org + +Why this exists +--------------- +Until now a /join submission went straight to the Google Sheet, with a raw +append to /data/leads.jsonl as a fire-and-forget side effect wrapped in a bare +except. The sheet is the only queryable copy, it has no record IDs, no status, +no way to mark a lead handled, and a failed write survives only as an ntfy push. + +This module makes a local SQLite database the source of truth. The sheet and +the ntfy push become *mirrors* of a record that already exists locally, and +each mirror's success or failure is recorded per record so it can be replayed. + +Designed for a future admin panel +--------------------------------- +The table is deliberately NOT a leads table. It is a generic record store keyed +by `kind`, so RSVPs, volunteer signups, popcorn orders or anything else the +admin panel eventually covers land in the same place with the same workflow +columns (status / assigned_to / notes) and the same audit trail. Adding a new +form means picking a new `kind` string, not migrating a schema. + +Stdlib only - sqlite3 ships with Python, so this adds no image dependencies. +""" + +import json +import os +import sqlite3 +import uuid +import datetime +from pathlib import Path + +DB_PATH = Path(os.environ.get("STORE_DB", "/data/scout73.db")) + +# Workflow states an admin panel may set. 'new' is the only one this app writes. +STATUSES = ("new", "contacted", "joined", "declined", "duplicate", "spam") + +# Mirror targets - external destinations a record is copied out to. +TARGETS = ("google_sheet", "ntfy") + +SCHEMA = """ +PRAGMA journal_mode=WAL; + +CREATE TABLE IF NOT EXISTS records ( + id TEXT PRIMARY KEY, + kind TEXT NOT NULL, + source TEXT NOT NULL, + unit TEXT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + display_name TEXT, + email TEXT, + phone TEXT, + payload TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'new', + assigned_to TEXT, + notes TEXT +); +CREATE INDEX IF NOT EXISTS idx_records_kind_created ON records(kind, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_records_status ON records(status); +CREATE INDEX IF NOT EXISTS idx_records_email ON records(email); + +CREATE TABLE IF NOT EXISTS mirrors ( + record_id TEXT NOT NULL, + target TEXT NOT NULL, + state TEXT NOT NULL, + attempts INTEGER NOT NULL DEFAULT 0, + last_error TEXT, + last_attempt_at TEXT, + PRIMARY KEY (record_id, target) +); +CREATE INDEX IF NOT EXISTS idx_mirrors_state ON mirrors(target, state); + +CREATE TABLE IF NOT EXISTS audit ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + record_id TEXT, + at TEXT NOT NULL, + actor TEXT NOT NULL, + action TEXT NOT NULL, + detail TEXT +); +CREATE INDEX IF NOT EXISTS idx_audit_record ON audit(record_id, id); + +CREATE TABLE IF NOT EXISTS meta ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL +); +""" + + +def _now(): + return datetime.datetime.now(datetime.timezone.utc).isoformat(timespec="seconds") + + +def _norm(ts): + """Normalise any incoming timestamp to UTC ISO-8601 with an offset. + + app.py builds its record timestamp with a naive datetime.now(), i.e. local + Eastern time. Storing that next to UTC audit rows would silently break both + sorting and `since=` filters in the admin panel, so everything that lands in + the DB is converted here. + """ + if not ts: + return _now() + try: + dt = datetime.datetime.fromisoformat(ts) + except Exception: + return _now() + if dt.tzinfo is None: + dt = dt.astimezone() # interpret as this host's local time + return dt.astimezone(datetime.timezone.utc).isoformat(timespec="seconds") + + +def connect(): + DB_PATH.parent.mkdir(parents=True, exist_ok=True) + con = sqlite3.connect(DB_PATH, timeout=10) + con.row_factory = sqlite3.Row + con.execute("PRAGMA foreign_keys=ON") + return con + + +def init(backfill_jsonl=None): + """Create the schema. Safe to call on every boot.""" + con = connect() + try: + con.executescript(SCHEMA) + con.commit() + if backfill_jsonl: + _backfill(con, Path(backfill_jsonl)) + finally: + con.close() + + +# --------------------------------------------------------------------------- +# Writes +# --------------------------------------------------------------------------- + +def insert_record(kind, source, payload, unit=None, display_name=None, + email=None, phone=None, created_at=None, record_id=None): + """Insert a record and return its id. Raises on failure - callers decide.""" + rid = record_id or str(uuid.uuid4()) + ts = _norm(created_at) + con = connect() + try: + con.execute( + "INSERT INTO records (id, kind, source, unit, created_at, updated_at," + " display_name, email, phone, payload, status)" + " VALUES (?,?,?,?,?,?,?,?,?,?, 'new')", + (rid, kind, source, unit, ts, ts, display_name, email, phone, + json.dumps(payload, ensure_ascii=False)), + ) + con.execute( + "INSERT INTO audit (record_id, at, actor, action, detail) VALUES (?,?,?,?,?)", + (rid, _now(), "system", "created", kind), + ) + con.commit() + finally: + con.close() + return rid + + +def set_mirror(record_id, target, ok, error=None): + """Record the outcome of an attempt to copy a record to an external target.""" + ts = _now() + state = "ok" if ok else "failed" + con = connect() + try: + con.execute( + "INSERT INTO mirrors (record_id, target, state, attempts, last_error, last_attempt_at)" + " VALUES (?,?,?,1,?,?)" + " ON CONFLICT(record_id, target) DO UPDATE SET" + " state=excluded.state," + " attempts=mirrors.attempts+1," + " last_error=excluded.last_error," + " last_attempt_at=excluded.last_attempt_at", + (record_id, target, state, (str(error)[:500] if error else None), ts), + ) + con.commit() + finally: + con.close() + + +def update_record(record_id, actor="admin", status=None, assigned_to=None, notes=None): + """Workflow update from the admin panel. Returns the updated row, or None.""" + sets, vals = [], [] + if status is not None: + if status not in STATUSES: + raise ValueError("unknown status: %s" % status) + sets.append("status=?"); vals.append(status) + if assigned_to is not None: + sets.append("assigned_to=?"); vals.append(assigned_to or None) + if notes is not None: + sets.append("notes=?"); vals.append(notes or None) + if not sets: + return get_record(record_id) + ts = _now() + sets.append("updated_at=?"); vals.append(ts) + vals.append(record_id) + con = connect() + try: + cur = con.execute("UPDATE records SET %s WHERE id=?" % ", ".join(sets), vals) + if cur.rowcount == 0: + return None + detail = json.dumps({k: v for k, v in + (("status", status), ("assigned_to", assigned_to), ("notes", notes)) + if v is not None}, ensure_ascii=False) + con.execute( + "INSERT INTO audit (record_id, at, actor, action, detail) VALUES (?,?,?,?,?)", + (record_id, ts, actor, "updated", detail), + ) + con.commit() + finally: + con.close() + return get_record(record_id) + + +def log(record_id, actor, action, detail=None): + con = connect() + try: + con.execute( + "INSERT INTO audit (record_id, at, actor, action, detail) VALUES (?,?,?,?,?)", + (record_id, _now(), actor, action, detail), + ) + con.commit() + finally: + con.close() + + +# --------------------------------------------------------------------------- +# Reads +# --------------------------------------------------------------------------- + +def _row(r): + d = dict(r) + if "payload" in d and d["payload"]: + try: + d["payload"] = json.loads(d["payload"]) + except Exception: + pass + return d + + +def get_record(record_id, with_audit=False): + con = connect() + try: + r = con.execute("SELECT * FROM records WHERE id=?", (record_id,)).fetchone() + if not r: + return None + out = _row(r) + out["mirrors"] = [dict(m) for m in con.execute( + "SELECT target, state, attempts, last_error, last_attempt_at" + " FROM mirrors WHERE record_id=?", (record_id,))] + if with_audit: + out["audit"] = [dict(a) for a in con.execute( + "SELECT at, actor, action, detail FROM audit WHERE record_id=? ORDER BY id", + (record_id,))] + return out + finally: + con.close() + + +def list_records(kind=None, status=None, since=None, q=None, limit=100, offset=0): + where, vals = [], [] + if kind: + where.append("kind=?"); vals.append(kind) + if status: + where.append("status=?"); vals.append(status) + if since: + where.append("created_at>=?"); vals.append(since) + if q: + where.append("(display_name LIKE ? OR email LIKE ? OR payload LIKE ?)") + vals += ["%%%s%%" % q] * 3 + sql = "SELECT * FROM records" + if where: + sql += " WHERE " + " AND ".join(where) + sql += " ORDER BY created_at DESC LIMIT ? OFFSET ?" + vals += [max(1, min(int(limit), 500)), max(0, int(offset))] + con = connect() + try: + rows = [_row(r) for r in con.execute(sql, vals)] + ids = [r["id"] for r in rows] + if ids: + marks = ",".join("?" * len(ids)) + mir = {} + for m in con.execute( + "SELECT record_id, target, state FROM mirrors WHERE record_id IN (%s)" % marks, ids): + mir.setdefault(m["record_id"], {})[m["target"]] = m["state"] + for r in rows: + r["mirrors"] = mir.get(r["id"], {}) + return rows + finally: + con.close() + + +def summary(): + """Counts for a dashboard: by kind, by status, and failed mirrors.""" + con = connect() + try: + return { + "total": con.execute("SELECT COUNT(*) c FROM records").fetchone()["c"], + "by_kind": {r["kind"]: r["c"] for r in con.execute( + "SELECT kind, COUNT(*) c FROM records GROUP BY kind")}, + "by_status": {r["status"]: r["c"] for r in con.execute( + "SELECT status, COUNT(*) c FROM records GROUP BY status")}, + "failed_mirrors": {r["target"]: r["c"] for r in con.execute( + "SELECT target, COUNT(*) c FROM mirrors WHERE state='failed' GROUP BY target")}, + } + finally: + con.close() + + +def failed_mirror_records(target, limit=50): + con = connect() + try: + return [_row(r) for r in con.execute( + "SELECT r.* FROM records r JOIN mirrors m ON m.record_id=r.id" + " WHERE m.target=? AND m.state='failed' ORDER BY r.created_at LIMIT ?", + (target, limit))] + finally: + con.close() + + +# --------------------------------------------------------------------------- +# One-time backfill of the pre-existing raw log +# --------------------------------------------------------------------------- + +def _backfill(con, path): + """Import /data/leads.jsonl once, so history is not stranded outside the DB.""" + done = con.execute("SELECT value FROM meta WHERE key='backfill_leads_jsonl'").fetchone() + if done or not path.exists(): + return + n = 0 + with open(path, encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + rec = json.loads(line) + except Exception: + continue + rid = str(uuid.uuid4()) + ts = _norm(rec.get("ts")) + con.execute( + "INSERT INTO records (id, kind, source, unit, created_at, updated_at," + " display_name, email, phone, payload, status)" + " VALUES (?,?,?,?,?,?,?,?,?,?, 'new')", + (rid, "join_lead", rec.get("source") or "greenlanescouts73.org", + unit_of(rec.get("interested_in", "")), ts, ts, + rec.get("parent_name"), rec.get("email"), rec.get("phone"), + json.dumps(rec, ensure_ascii=False)), + ) + con.execute( + "INSERT INTO audit (record_id, at, actor, action, detail) VALUES (?,?,?,?,?)", + (rid, _now(), "system", "backfilled", "leads.jsonl"), + ) + n += 1 + con.execute("INSERT INTO meta (key, value) VALUES ('backfill_leads_jsonl', ?)", + ("%s rows at %s" % (n, _now()),)) + con.commit() + print("store: backfilled %s rows from %s" % (n, path), flush=True) + + +def unit_of(interested_in): + """Best-effort unit tag so the admin panel can filter pack vs troop.""" + s = (interested_in or "").lower() + pack = "pack" in s or "cub" in s + troop = "troop" in s or "scouts bsa" in s + if pack and troop or "both" in s: + return "both" + if pack: + return "pack" + if troop: + return "troop" + return None diff --git a/docker-compose.yml b/docker-compose.yml index 4cb6cab..58382e6 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -24,6 +24,7 @@ services: - NTFY_BASE=${NTFY_BASE} - NTFY_TOPIC=${NTFY_TOPIC} - NTFY_TOKEN=${NTFY_TOKEN} + - ADMIN_TOKEN=${ADMIN_TOKEN} volumes: - /srv/scout-website/data:/data - /srv/scout-website/secrets/service_account.json:/app/service_account.json:ro