Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,13 @@ jobs:
- name: Install dependencies
run: |
python -m pip install --upgrade pip
pip install -r requirements.txt ruff
pip install -r requirements.txt ruff==0.15.20

- name: Lint
run: ruff check app.py prizm_client.py cache_manager_new.py cache_cli.py test_prizm_api.py
run: ruff check app.py prizm_client.py cache_manager_new.py cache_cli.py cohort_queue.py historical_import.py test_prizm_api.py test_cohort_queue.py test_historical_import.py

- name: Test
run: python -m unittest test_prizm_api.py
run: python -m unittest test_prizm_api.py test_cohort_queue.py test_historical_import.py

node:
name: Node tests and build
Expand Down
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -56,3 +56,5 @@ debug_screenshots/
missing_postal_codes.txt
postal_code_export*.csv
postal_code_v2_export*.csv

*.migration.lock
17 changes: 15 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ GET /health
GET /dashboard
```

The dashboard shows cache totals, recent daily additions, recent failures, searchable cached postal-code data, CSV export, and a manual weekly-report send button. It uses HTTP Basic Auth.
The root URL and `/dashboard` show cache totals, daily lookup outcomes, failure history (including uncached quota/network errors), searchable and paginated postal-code records with expandable details, complete cache/failure CSV exports, and weekly-report preview. It uses HTTP Basic Auth and fails closed without a configured password or API key. Dashboard credentials provide read-only reporting access; sending an email requires the API key.

Recommended dashboard/reporting variables:

Expand Down Expand Up @@ -74,7 +74,9 @@ Useful if you want the full segment reference data without postal-code lookup.

```http
GET /api/cache/entries?status=success&search=V8A&limit=500
GET /api/cache/export.csv
GET /api/cache/export.csv?include_expired=1
GET /api/lookups?failures=1&days=7&limit=100&offset=0
GET /api/lookups/failures.csv?days=3650
```

### Weekly Report
Expand Down Expand Up @@ -192,3 +194,14 @@ python3 -m unittest test_prizm_api.py
- Urban postal codes still depend on the Environics geocoder API. During triage on June 21, 2026, that endpoint returned `403 Quota Exceeded`, so urban postal-code lookups may fail until the public quota replenishes or a licensed/geocoding credential is available.
- `prizm_cache_v2.db`, CSV exports, debug screenshots, and `node_modules` are runtime/local artifacts and are ignored by Docker.
- The old crash was most likely caused by a globally reused Selenium Chrome session combined with Flask debug mode and repeated screenshot/page-source capture. Removing the browser from request handling eliminates that failure path.

## Operational behavior

- Weekly metrics use a trailing seven-day window with UTC daily dates. Lookup success means an API result, not a successful Salesforce update. Cache hits, upstream attempts, upstream successes/failures, and unique postal codes are separate metrics. First captures mean the first successful upstream result in retained lookup-event history; imported/earlier records may predate that history.
- The daily Salesforce selector no longer retries an assigned household solely because net worth is unavailable. It preserves an existing net-worth value when no historical reference exists. Retryable quota/network errors leave the household eligible for a later attempt.
- Salesforce household updates now roll back the batch and surface a failed job if any update fails. This exposes errors previously hidden by ignored partial-update results; inspect the job error before retrying a persistently failing batch.
- The current public reference data includes dwelling type but no net worth. Net worth shown here is historical segment reference data of mixed vintage, not a current postal-code-specific estimate. Missing values remain unavailable. Some dwelling types cannot be represented by the existing Salesforce picklist and map to `None or Unverifiable`.
- `MAX_BATCH_POSTAL_CODES=10` limits this application's request size; it does not establish an upstream allowance. A claimed three-per-day restriction remains unverified. Confirm licensed access with Environics before adjusting collection volume.
- Weekly delivery already uses Salesforce's scheduled sender; do not create a second Railway email schedule. Salesforce job completion is not an inbox-delivery receipt.

See [September 6 investigation and upgrade scope](docs/2026-09-06-upgrade.md).
191 changes: 167 additions & 24 deletions app.py

Large diffs are not rendered by default.

120 changes: 113 additions & 7 deletions cache_manager_new.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
"""

import csv
import fcntl
import io
import json
import logging
Expand All @@ -16,6 +17,8 @@
from datetime import datetime, timedelta
from typing import Any, Dict, List, Optional

from segment_net_worth import average_household_net_worth, average_household_net_worth_amount

logger = logging.getLogger(__name__)


Expand All @@ -25,7 +28,23 @@ class CacheManager:
def __init__(self, db_path: str = None, cache_duration_days: int = None):
self.db_path = db_path or os.environ.get("PRIZM_CACHE_DB_PATH", "prizm_cache_v2.db")
self.cache_duration_days = cache_duration_days or int(os.environ.get("PRIZM_CACHE_DURATION_DAYS", "90"))
self._init_database()
# Serialize startup migrations across Gunicorn workers. Keep the original
# production database before this release adds tables/columns/triggers.
with open(self.db_path + ".migration.lock", "a") as lock:
fcntl.flock(lock, fcntl.LOCK_EX)
if os.environ.get("RAILWAY_ENVIRONMENT_ID") and os.path.exists(self.db_path):
snapshot = self.db_path + ".before-cohort-upgrade.db"
if not os.path.exists(snapshot):
source = sqlite3.connect(self.db_path)
target = sqlite3.connect(snapshot + ".tmp")
try:
source.backup(target)
finally:
target.close()
source.close()
os.replace(snapshot + ".tmp", snapshot)
logger.info("Saved pre-upgrade database snapshot")
self._init_database()

@contextmanager
def _connect(self):
Expand Down Expand Up @@ -218,6 +237,7 @@ def _init_database(self):
cursor,
"postal_code_cache",
{
"historical_source": "TEXT",
"segment_number": "TEXT",
"segment_name": "TEXT",
"segment_description": "TEXT",
Expand Down Expand Up @@ -281,6 +301,27 @@ def _init_database(self):
cursor.execute("CREATE INDEX IF NOT EXISTS idx_lookup_events_source ON lookup_events (source)")
cursor.execute("CREATE INDEX IF NOT EXISTS idx_lookup_events_postal_code ON lookup_events (postal_code)")

cursor.executescript("""
CREATE TABLE IF NOT EXISTS capture_history (
postal_code TEXT PRIMARY KEY, captured_at TEXT NOT NULL
);
INSERT OR IGNORE INTO capture_history
SELECT postal_code, COALESCE(cached_at, CURRENT_TIMESTAMP)
FROM postal_code_cache WHERE status='success';
CREATE TRIGGER IF NOT EXISTS remember_prizm_capture_insert
AFTER INSERT ON postal_code_cache WHEN NEW.status='success'
BEGIN
INSERT OR IGNORE INTO capture_history VALUES
(NEW.postal_code, COALESCE(NEW.cached_at, CURRENT_TIMESTAMP));
END;
CREATE TRIGGER IF NOT EXISTS remember_prizm_capture_update
AFTER UPDATE ON postal_code_cache WHEN NEW.status='success'
BEGIN
INSERT OR IGNORE INTO capture_history VALUES
(NEW.postal_code, COALESCE(NEW.cached_at, CURRENT_TIMESTAMP));
END;
""")

conn.commit()
logger.info("Cache database initialized at %s", self.db_path)

Expand Down Expand Up @@ -682,6 +723,15 @@ def list_cache_entries(

def _dashboard_row(self, row: sqlite3.Row) -> Dict[str, Any]:
data = self._row_to_response_dict(row)
amount = average_household_net_worth_amount(data.get("segment_number"))
if amount is not None:
if not data.get("average_household_net_worth"):
data["average_household_net_worth"] = average_household_net_worth(data.get("segment_number"))
if data.get("average_household_net_worth_amount") is None:
data["average_household_net_worth_amount"] = amount
data["net_worth_source"] = "Historical segment reference" if data.get("average_household_net_worth") else "Unavailable"
data["data_source"] = ("Historical import: " + row["historical_source"]
if row["historical_source"] and row["cached_at"] is None else "Current dataset")
data.update(
{
"cached_at": row["cached_at"],
Expand Down Expand Up @@ -723,7 +773,15 @@ def get_lookup_event_summary(self, days: int = 7) -> Dict[str, Any]:
SUM(CASE WHEN status = 'success' THEN 1 ELSE 0 END) successful,
SUM(CASE WHEN status != 'success' THEN 1 ELSE 0 END) failed,
SUM(CASE WHEN from_cache = 1 THEN 1 ELSE 0 END) cache_hits,
SUM(CASE WHEN source = 'upstream' THEN 1 ELSE 0 END) upstream_attempts
SUM(CASE WHEN source = 'upstream' THEN 1 ELSE 0 END) upstream_attempts,
COUNT(DISTINCT postal_code) unique_postal_codes,
SUM(CASE WHEN source = 'upstream' AND status = 'success' THEN 1 ELSE 0 END) upstream_successful,
SUM(CASE WHEN source = 'upstream' AND status != 'success' THEN 1 ELSE 0 END) upstream_failed,
SUM(CASE WHEN source = 'upstream' AND status = 'success' AND NOT EXISTS (
SELECT 1 FROM lookup_events prior
WHERE prior.postal_code = lookup_events.postal_code
AND prior.status = 'success' AND prior.id < lookup_events.id
) THEN 1 ELSE 0 END) newly_captured
FROM lookup_events
WHERE requested_at >= datetime('now', ?)
""",
Expand All @@ -738,7 +796,15 @@ def get_lookup_event_summary(self, days: int = 7) -> Dict[str, Any]:
SUM(CASE WHEN status = 'success' THEN 1 ELSE 0 END) successful,
SUM(CASE WHEN status != 'success' THEN 1 ELSE 0 END) failed,
SUM(CASE WHEN from_cache = 1 THEN 1 ELSE 0 END) cache_hits,
SUM(CASE WHEN source = 'upstream' THEN 1 ELSE 0 END) upstream_attempts
SUM(CASE WHEN source = 'upstream' THEN 1 ELSE 0 END) upstream_attempts,
COUNT(DISTINCT postal_code) unique_postal_codes,
SUM(CASE WHEN source = 'upstream' AND status = 'success' THEN 1 ELSE 0 END) upstream_successful,
SUM(CASE WHEN source = 'upstream' AND status != 'success' THEN 1 ELSE 0 END) upstream_failed,
SUM(CASE WHEN source = 'upstream' AND status = 'success' AND NOT EXISTS (
SELECT 1 FROM lookup_events prior
WHERE prior.postal_code = lookup_events.postal_code
AND prior.status = 'success' AND prior.id < lookup_events.id
) THEN 1 ELSE 0 END) newly_captured
FROM lookup_events
WHERE requested_at >= datetime('now', ?)
GROUP BY date(requested_at)
Expand All @@ -757,11 +823,43 @@ def get_dashboard_summary(self) -> Dict[str, Any]:
"cache_stats": self.get_cache_stats(),
"daily_cache_counts": self.get_daily_cache_counts(30),
"lookup_events_7d": self.get_lookup_event_summary(7),
"recent_failures": self.list_cache_entries(status="error", limit=50),
"recent_failures": self.list_lookup_events(failures_only=True, days=7, limit=50),
}

def list_lookup_events(self, failures_only=False, days=7, limit=500, offset=0):
conditions = ["requested_at >= datetime('now', ?)"]
if failures_only:
conditions.append("status != 'success'")
with self._connect() as conn:
rows = conn.execute(
f"SELECT * FROM lookup_events WHERE {' AND '.join(conditions)} "
"ORDER BY requested_at DESC, id DESC LIMIT ? OFFSET ?",
(f"-{max(1, min(int(days), 3650))} days", max(1, min(int(limit), 5000)), max(0, int(offset))),
).fetchall()
return [dict(row) for row in rows]

@staticmethod
def _csv_value(value):
if isinstance(value, str) and value.lstrip().startswith(('=', '+', '-', '@')):
return "'" + value
return value

def export_failures_csv(self, days=7):
output = io.StringIO()
fields = ['postal_code', 'requested_at', 'status', 'source', 'message', 'endpoint', 'batch_id']
writer = csv.DictWriter(output, fieldnames=fields)
writer.writeheader()
offset = 0
while True:
rows = self.list_lookup_events(failures_only=True, days=days, limit=5000, offset=offset)
for row in rows:
writer.writerow({field: self._csv_value(row.get(field, '')) for field in fields})
if len(rows) < 5000:
break
offset += len(rows)
return output.getvalue()

def export_cache_csv(self, include_expired: bool = False) -> str:
rows = self.list_cache_entries(limit=5000, include_expired=include_expired)
output = io.StringIO()
fieldnames = [
"postal_code",
Expand All @@ -774,6 +872,8 @@ def export_cache_csv(self, include_expired: bool = False) -> str:
"average_household_income",
"average_household_net_worth",
"average_household_net_worth_amount",
"net_worth_source",
"data_source",
"education",
"urbanity",
"occupation",
Expand All @@ -796,8 +896,14 @@ def export_cache_csv(self, include_expired: bool = False) -> str:
]
writer = csv.DictWriter(output, fieldnames=fieldnames)
writer.writeheader()
for row in rows:
writer.writerow({field: row.get(field, "") for field in fieldnames})
offset = 0
while True:
rows = self.list_cache_entries(limit=5000, offset=offset, include_expired=include_expired)
for row in rows:
writer.writerow({field: self._csv_value(row.get(field, "")) for field in fieldnames})
if len(rows) < 5000:
break
offset += len(rows)
return output.getvalue()

def clear_cache(self) -> bool:
Expand Down
Loading
Loading