Two real leads on 2026-08-26 landed in the DB with no Sheet row and no mirror
row at all: the empty-env window meant the Sheet call never ran. Because
/mirrors/failed looked for state='failed', a mirror that was never attempted
read exactly like one that was never owed, and the summary reported a clean
failed_mirrors: {} while two families were missing from the sheet.
/join now writes a `pending` mirror row for every target BEFORE attempting any
of them, so the gap between "owed" and "done" is a row rather than an absence.
set_mirror_pending never downgrades an attempted mirror. failed_mirror_records
covers failed and pending, since both need replaying, and summary reports
pending_mirrors alongside failed_mirrors.
sheet_append also stamped the Sheet with datetime.now(), so replaying a lead
filed it under the day of the replay rather than the day the family submitted.
It now takes the submission timestamp off the record and only falls back to now
when there isn't one.
339 lines
12 KiB
Python
339 lines
12 KiB
Python
"""
|
|
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)
|