""" 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") # Announcement guardrails. Enforced at write time so a bad notice never # reaches a render. Both are rejections, never silent truncation: clipping # someone's cancellation mid-sentence is worse than making them shorten it. ANNOUNCEMENT_MAX_CHARS = 200 ANNOUNCEMENT_MAX_LIVE = 3 ANNOUNCEMENT_LEVELS = ("info", "urgent") 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 ); CREATE TABLE IF NOT EXISTS announcements ( id TEXT PRIMARY KEY, created_at TEXT NOT NULL, message TEXT NOT NULL, level TEXT NOT NULL DEFAULT 'info', starts_at TEXT NOT NULL, ends_at TEXT NOT NULL, link_url TEXT, link_text TEXT, created_by TEXT, revoked_at TEXT ); CREATE INDEX IF NOT EXISTS idx_announcements_window ON announcements(revoked_at, starts_at, ends_at); -- Other Scouting units in the Continental District, published as a courtesy -- directory. This is somebody else's data: it comes from a hand-typed district -- document and goes stale without telling us, so `verified_at` is shown to the -- reader rather than kept as bookkeeping. Pack 73 and Troop 73 are deliberately -- NOT rows here - the rest of this site is us. CREATE TABLE IF NOT EXISTS nearby_units ( id TEXT PRIMARY KEY, unit_type TEXT NOT NULL, unit_number TEXT NOT NULL, serves TEXT, chartered_org TEXT, street TEXT, town TEXT, area TEXT, meets TEXT, notes TEXT, link_url TEXT, contact TEXT, sort_order INTEGER NOT NULL DEFAULT 100, active INTEGER NOT NULL DEFAULT 1, source TEXT, verified_at TEXT, updated_at TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_nearby_units_listing ON nearby_units(active, sort_order, unit_number); """ 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) # ---------------------------------------------------------------------------- # Announcements # ---------------------------------------------------------------------------- # # A site-wide notice with a start and an end. It exists because "tonight's # meeting is cancelled, the lot is flooded" at 4pm on a Tuesday is the one # string on this site where a deploy is the wrong latency. # # ends_at is REQUIRED. That is the whole point: nothing has to be remembered # and taken down. An announcement with no end is site copy, and site copy # lives in git where it has a diff. # # Nothing is ever hard-deleted. Taking one down early sets revoked_at, so the # record of what the site said, and when, survives. class AnnouncementRejected(Exception): """Raised when a write breaks a guardrail. Carries the HTTP status the admin API should return, so the caps live here rather than in the route.""" def __init__(self, status, detail, extra=None): super().__init__(detail) self.status = status self.detail = detail self.extra = extra or {} def _live_at(con, when): return con.execute( "SELECT * FROM announcements" " WHERE revoked_at IS NULL AND starts_at <= ? AND ends_at > ?" " ORDER BY CASE level WHEN 'urgent' THEN 0 ELSE 1 END, starts_at DESC", (when, when), ).fetchall() def active_announcements(now=None): """The notices that should render, best first. Urgent outranks info, then most recent. Recency alone would let a routine Wednesday notice bury a Tuesday cancellation that is still live. Capped at ANNOUNCEMENT_MAX_LIVE as a floor under the render even if rows got in past the write check - the banner is never allowed to be unbounded. """ when = now or _now() con = connect() try: return [dict(r) for r in _live_at(con, when)[:ANNOUNCEMENT_MAX_LIVE]] finally: con.close() def list_announcements(include_expired=False, limit=100): """Every announcement with a computed live/expired/revoked state. The state is returned rather than inferred, so 'why is my notice not showing' is answerable from the API instead of from the homepage. """ now = _now() con = connect() try: sql = "SELECT * FROM announcements" if not include_expired: sql += " WHERE revoked_at IS NULL AND ends_at > '%s'" % now sql += " ORDER BY starts_at DESC LIMIT ?" rows = [dict(r) for r in con.execute(sql, (limit,)).fetchall()] finally: con.close() live_ids = {r["id"] for r in active_announcements(now)} for r in rows: if r["revoked_at"]: r["state"] = "revoked" elif r["ends_at"] <= now: r["state"] = "expired" elif r["starts_at"] > now: r["state"] = "scheduled" elif r["id"] in live_ids: r["state"] = "live" else: r["state"] = "over_cap" return rows def create_announcement(message, ends_at, starts_at=None, level="info", link_url=None, link_text=None, created_by=None): message = (message or "").strip() if not message: raise AnnouncementRejected(422, "message is required") if len(message) > ANNOUNCEMENT_MAX_CHARS: raise AnnouncementRejected(422, ( "message is %d characters and the cap is %d. Put the long version on a " "documents page and link to it with link_url." % (len(message), ANNOUNCEMENT_MAX_CHARS))) if level not in ANNOUNCEMENT_LEVELS: raise AnnouncementRejected(422, "level must be one of %s" % (ANNOUNCEMENT_LEVELS,)) if not ends_at: raise AnnouncementRejected(422, ( "ends_at is required. An announcement that never expires is site copy, " "and site copy belongs in the repo where it has a diff.")) if link_text and not link_url: raise AnnouncementRejected(422, "link_text without link_url has nothing to point at") if link_url and not str(link_url).startswith(("/", "https://")): raise AnnouncementRejected(422, "link_url must be site-relative or https") starts = _norm(starts_at) if starts_at else _now() ends = _norm(ends_at) if ends <= starts: raise AnnouncementRejected(422, "ends_at must be after starts_at") aid = str(uuid.uuid4()) con = connect() try: live = _live_at(con, starts) if len(live) >= ANNOUNCEMENT_MAX_LIVE: raise AnnouncementRejected(409, ( "%d announcements are already live at that start time and the cap is %d. " "Revoke one first." % (len(live), ANNOUNCEMENT_MAX_LIVE)), {"live": [{k: r[k] for k in ("id", "message", "level", "ends_at")} for r in live]}) con.execute( "INSERT INTO announcements (id, created_at, message, level, starts_at," " ends_at, link_url, link_text, created_by, revoked_at)" " VALUES (?,?,?,?,?,?,?,?,?,NULL)", (aid, _now(), message, level, starts, ends, link_url or None, link_text or None, created_by or None)) con.commit() finally: con.close() return get_announcement(aid) def get_announcement(aid): con = connect() try: row = con.execute("SELECT * FROM announcements WHERE id=?", (aid,)).fetchone() return dict(row) if row else None finally: con.close() def revoke_announcement(aid): """Take one down early. Never deletes - the site's history is the point.""" con = connect() try: cur = con.execute( "UPDATE announcements SET revoked_at=? WHERE id=? AND revoked_at IS NULL", (_now(), aid)) con.commit() return cur.rowcount > 0 finally: con.close()