""" store.py - internal record store for greenlanescouts73.org What this holds --------------- One table per thing the site collects. Today that is `join_leads`: the /join interest form, one row per submission, one column per form field. A /join submission lands here FIRST. The Google Sheet and the ntfy push are mirrors of a row that already exists locally, and each mirror's outcome is recorded per row so a failure is replayable instead of just shouted. Why join_leads and not a generic table ------------------------------------- This started as a generic `records` table keyed by `kind`, carrying status/assigned_to/notes for a future admin panel. That was wrong twice over. The PII columns (display_name/email/phone) only make sense for a lead, and the workflow columns describe an outreach process that has not been designed yet - one that is plainly one-to-many, since a family gets contacted more than once. Guessing at it in three columns would have locked in the wrong shape. So: a table per form, columns for the form's own fields, and outreach gets its own table when the process is actually known. Anything the form adds later that is not worth a column still survives in `payload`, which holds the submission verbatim. 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")) # Mirror targets - external destinations a row is copied out to. TARGETS = ("google_sheet", "ntfy") SCHEMA = """ PRAGMA journal_mode=WAL; CREATE TABLE IF NOT EXISTS join_leads ( id TEXT PRIMARY KEY, submitted_at TEXT NOT NULL, recorded_at TEXT NOT NULL, parent_name TEXT NOT NULL, email TEXT NOT NULL, phone TEXT, interested_in TEXT, children TEXT, heard_from TEXT, heard_from_detail TEXT, message TEXT, payload TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_join_leads_submitted ON join_leads(submitted_at DESC); CREATE INDEX IF NOT EXISTS idx_join_leads_email ON join_leads(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 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 submission timestamp with a naive datetime.now(), i.e. local Eastern time. Storing that next to UTC rows would silently break both sorting and `since=` filters, 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_lead(parent_name, email, phone=None, interested_in=None, children=None, heard_from=None, heard_from_detail=None, message=None, payload=None, submitted_at=None, record_id=None): """Insert one /join submission and return its id. Raises on failure - callers decide what a failed local write means.""" rid = record_id or str(uuid.uuid4()) submitted = _norm(submitted_at) con = connect() try: con.execute( "INSERT INTO join_leads (id, submitted_at, recorded_at, parent_name, email," " phone, interested_in, children, heard_from, heard_from_detail, message, payload)" " VALUES (?,?,?,?,?,?,?,?,?,?,?,?)", (rid, submitted, _now(), parent_name, email, phone or None, interested_in or None, children or None, heard_from or None, heard_from_detail or None, message or None, json.dumps(payload or {}, ensure_ascii=False)), ) con.commit() finally: con.close() return rid def set_mirror_pending(record_id, target): """Register a mirror as owed, before it is attempted. Without this a mirror that never ran leaves no row at all, which reads identically to a mirror that was never owed. That is how two leads went missing from the Sheet on 2026-08-26 with a clean `failed_mirrors: {}`. Never downgrades an existing row - a mirror already ok or failed has been attempted, and its outcome stands. """ con = connect() try: con.execute( "INSERT INTO mirrors (record_id, target, state, attempts, last_attempt_at)" " VALUES (?,?,'pending',0,NULL)" " ON CONFLICT(record_id, target) DO NOTHING", (record_id, target), ) con.commit() finally: con.close() def set_mirror(record_id, target, ok, error=None): """Record the outcome of an attempt to copy a row 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() # ---------------------------------------------------------------------------- # Reads # ---------------------------------------------------------------------------- def _row(r): d = dict(r) if d.get("payload"): try: d["payload"] = json.loads(d["payload"]) except Exception: pass return d def get_lead(record_id): con = connect() try: r = con.execute("SELECT * FROM join_leads 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,))] return out finally: con.close() def list_leads(since=None, q=None, limit=100, offset=0): where, vals = [], [] if since: where.append("submitted_at>=?"); vals.append(since) if q: where.append("(parent_name LIKE ? OR email LIKE ? OR children LIKE ?" " OR heard_from LIKE ? OR message LIKE ?)") vals += ["%%%s%%" % q] * 5 sql = "SELECT * FROM join_leads" if where: sql += " WHERE " + " AND ".join(where) sql += " ORDER BY submitted_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: totals, interest split, referral sources, failed mirrors.""" con = connect() try: return { "total": con.execute("SELECT COUNT(*) c FROM join_leads").fetchone()["c"], "by_interest": {r["interested_in"]: r["c"] for r in con.execute( "SELECT interested_in, COUNT(*) c FROM join_leads GROUP BY interested_in")}, "by_heard_from": {r["heard_from"]: r["c"] for r in con.execute( "SELECT heard_from, COUNT(*) c FROM join_leads GROUP BY heard_from" " ORDER BY c DESC")}, "failed_mirrors": {r["target"]: r["c"] for r in con.execute( "SELECT target, COUNT(*) c FROM mirrors WHERE state='failed' GROUP BY target")}, "pending_mirrors": {r["target"]: r["c"] for r in con.execute( "SELECT target, COUNT(*) c FROM mirrors WHERE state='pending' GROUP BY target")}, } finally: con.close() def failed_mirror_records(target, limit=50): """Rows owing a copy to `target`: attempted and failed, or never attempted. `pending` is included deliberately. A mirror that never ran needs replaying just as much as one that ran and failed, and it is the case that hid two leads on 2026-08-26. """ con = connect() try: return [_row(r) for r in con.execute( "SELECT l.* FROM join_leads l JOIN mirrors m ON m.record_id=l.id" " WHERE m.target=? AND m.state IN ('failed','pending')" " ORDER BY l.submitted_at LIMIT ?", (target, limit))] finally: con.close() # ---------------------------------------------------------------------------- # One-time backfill of the pre-existing raw log # ---------------------------------------------------------------------------- def split_heard_from(source): """app.py used to fold source_place / source_other into one composed string ('Other: Chocolate booth'), destroying the raw answer. Recover both halves.""" s = (source or "").strip() if s.startswith("School/Daycare: "): return "Other school or daycare", s[len("School/Daycare: "):].strip() if s.startswith("Other: "): return "Other", s[len("Other: "):].strip() return s or None, None 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 heard, detail = split_heard_from(rec.get("source")) ts = _norm(rec.get("ts")) con.execute( "INSERT INTO join_leads (id, submitted_at, recorded_at, parent_name, email," " phone, interested_in, children, heard_from, heard_from_detail, message, payload)" " VALUES (?,?,?,?,?,?,?,?,?,?,?,?)", (str(uuid.uuid4()), ts, _now(), rec.get("parent_name") or "(unknown)", rec.get("email") or "", rec.get("phone") or None, rec.get("interested_in") or None, rec.get("children") or None, heard, detail, rec.get("message") or None, json.dumps(rec, ensure_ascii=False)), ) 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)