Files
scout-website/app/store.py
T

542 lines
20 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")
# 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()