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
This commit is contained in:
Mike Wichers
2026-08-26 19:38:33 -04:00
parent 5331e38d62
commit be2d81b3ab
4 changed files with 506 additions and 0 deletions
+375
View File
@@ -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