From b640fd683b49cb95495b765dda440b23429da57a Mon Sep 17 00:00:00 2001 From: heybeaux Date: Sun, 6 Sep 2026 15:58:18 -0700 Subject: [PATCH 1/5] fix: unblock PRIZM enrichment and distinguish fresh captures in reports --- README.md | 17 ++- app.py | 106 ++++++++++++++---- cache_manager_new.py | 75 ++++++++++++- docs/2026-09-06-upgrade.md | 33 ++++++ salesforce/classes/WCPrizmDailyEnrichment.cls | 16 ++- .../classes/WCPrizmDailyEnrichmentTest.cls | 56 +++++++++ .../classes/WCPrizmWeeklyReportEmail.cls | 5 +- test_prizm_api.py | 101 +++++++++++++++++ 8 files changed, 373 insertions(+), 36 deletions(-) create mode 100644 docs/2026-09-06-upgrade.md diff --git a/README.md b/README.md index bde5305..4bc43fd 100644 --- a/README.md +++ b/README.md @@ -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: @@ -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 @@ -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). diff --git a/app.py b/app.py index a40450f..11fc73d 100644 --- a/app.py +++ b/app.py @@ -57,7 +57,9 @@ .muted { color:var(--muted); } .section-title { display:flex; justify-content:space-between; align-items:end; gap:12px; margin:28px 0 12px; } .section-title h2 { margin:0; font-size:20px; } - pre { white-space:pre-wrap; } + pre { white-space:pre-wrap; overflow-wrap:anywhere; max-width:540px; } + #reportStatus { white-space:pre-wrap; } + table { display:block; overflow-x:auto; } @media (max-width:900px) { .grid { grid-template-columns:repeat(2,minmax(0,1fr)); } main, header { padding-left:18px; padding-right:18px; } table { font-size:13px; } } @media (max-width:580px) { .grid { grid-template-columns:1fr; } th:nth-child(5),td:nth-child(5),th:nth-child(6),td:nth-child(6){display:none;} } @@ -72,25 +74,27 @@
–
postal codes cached
–
successful cached codes
–
failed / unassigned cached codes
-
–
lookups recorded this week
+
–
lookups in the trailing 7 days
-

Export and reports

Download the full cache or send the configured weekly email report.
+

Export and reports

Download cached data and failed attempts, or preview the scheduled weekly report.
- Export cached data CSV + Export cached data CSV + Export failure history CSV
-

Recent daily cache additions

Based on when each postal code was first cached or refreshed.
+

Daily lookup outcomes

Trailing 7 days, grouped by UTC date. Cache hits do not represent new captures.
+

Loading…
-

Postal codes

Search cached successes, failures, and invalid records.
+

Postal codes

Search all cached records, including expired entries. Net worth is historical segment reference data; unavailable values remain blank.
@@ -98,12 +102,19 @@
- +
Postal codeStatusSegmentNameHome typeIncomeNet worthCachedMessage
Postal codeStatusSegmentNameHome typeIncomeHistorical net worthCachedMessage
Loading…
+
+

Failed lookup history

+

Includes quota/network errors that are not stored in the cache. Dates are UTC. Export downloads all retained failures.

+
Loading…
+
@@ -210,7 +238,7 @@ def has_valid_api_key() -> bool: def has_valid_dashboard_auth() -> bool: username, password = dashboard_credentials() if not password: - return True + return False provided_username, provided_password = parse_basic_auth() return secrets.compare_digest(provided_username or "", username or "") and secrets.compare_digest(provided_password or "", password) @@ -226,12 +254,15 @@ def require_authentication(): if request.path == "/health" or request.method == "OPTIONS": return None - dashboard_paths = {"/", "/dashboard", "/api/dashboard/summary", "/api/cache/entries", "/api/cache/export.csv", "/api/reports/weekly", "/api/reports/weekly/send"} + dashboard_paths = {"/", "/dashboard", "/api/dashboard/summary", "/api/cache/entries", "/api/cache/export.csv", "/api/reports/weekly", "/api/lookups", "/api/lookups/failures.csv"} if request.path in dashboard_paths: - if has_valid_dashboard_auth() or has_valid_api_key(): + if has_valid_dashboard_auth() or (os.environ.get("PRIZM_API_KEY") and has_valid_api_key()): return None return dashboard_auth_required() + if request.path == "/api/reports/weekly/send" and not os.environ.get("PRIZM_API_KEY"): + return jsonify({"error": "Email sending requires a configured API key"}), 503 + if request.path.startswith("/api/") and has_valid_api_key(): return None @@ -304,6 +335,7 @@ def get_prizm_code(postal_code: str, endpoint: str = "single", batch_id: Optiona "home_type": "", "status": "error", "message": str(exc), + "retryable": True, } if should_cache: @@ -393,6 +425,24 @@ def dashboard_summary(): return jsonify(cache_manager.get_dashboard_summary()) +@app.route("/api/lookups", methods=["GET"]) +def lookup_events(): + rows = cache_manager.list_lookup_events( + failures_only=request.args.get("failures") == "1", + days=request.args.get("days", 7, type=int), + limit=request.args.get("limit", 100, type=int), + offset=request.args.get("offset", 0, type=int), + ) + return jsonify({"entries": rows, "count": len(rows)}) + + +@app.route("/api/lookups/failures.csv", methods=["GET"]) +def export_lookup_failures(): + response = Response(cache_manager.export_failures_csv(request.args.get("days", 7, type=int)), mimetype="text/csv") + response.headers["Content-Disposition"] = "attachment; filename=prizm-failures.csv" + return response + + @app.route("/api/cache/entries", methods=["GET"]) def get_cache_entries(): entries = cache_manager.list_cache_entries( @@ -483,34 +533,44 @@ def get_debug_html(postal_code): def build_weekly_report(days: int = 7) -> Dict[str, Any]: summary = cache_manager.get_dashboard_summary() - events = summary.get("lookup_events_7d") or cache_manager.get_lookup_event_summary(days) + events = cache_manager.get_lookup_event_summary(days) + summary["lookup_events_7d"] = events stats = summary.get("cache_stats", {}) - daily_counts = summary.get("daily_cache_counts", [])[:days] - failures = summary.get("recent_failures", [])[:20] + daily_counts = events.get("by_day", []) + failures = cache_manager.list_lookup_events(failures_only=True, days=days, limit=20) subject = os.environ.get("WEEKLY_REPORT_SUBJECT", "PRIZM weekly report") lines = [ "PRIZM weekly report", "", + f"Window: trailing {days} days; daily dates are UTC.", + "API success does not prove Salesforce household updates or email delivery.", f"Lookups recorded: {events.get('lookups', 0)}", f"Successful lookups: {events.get('successful', 0)}", f"Failed lookups: {events.get('failed', 0)}", f"Cache hits: {events.get('cache_hits', 0)}", f"Upstream attempts: {events.get('upstream_attempts', 0)}", + f"Upstream successful: {events.get('upstream_successful', 0)}", + f"Upstream failed: {events.get('upstream_failed', 0)}", + f"Unique postal codes requested: {events.get('unique_postal_codes', 0)}", + f"New captures (first successful lookup in recorded history): {events.get('newly_captured', 0)}", + "Historical net worth is a segment reference, not current postal-code wealth.", "", f"Total active cached postal codes: {stats.get('valid_entries', 0)}", f"Cached successes: {(stats.get('status_breakdown') or {}).get('success', 0)}", f"Cached failed/unassigned: {((stats.get('status_breakdown') or {}).get('error', 0) + (stats.get('status_breakdown') or {}).get('invalid', 0))}", "", - "Daily cache additions:", + "Daily lookup outcomes (UTC):", ] if daily_counts: for row in daily_counts: - lines.append(f"- {row.get('day')}: {row.get('total', 0)} total, {row.get('successful', 0)} success, {row.get('failed', 0)} failed") + lines.append(f"- {row.get('day')}: {row.get('lookups', 0)} lookups, {row.get('cache_hits', 0)} cached, {row.get('upstream_successful', 0)} upstream successes, {row.get('failed', 0)} failures") else: lines.append("- None recorded") - lines.extend(["", "Recent failures/unassigned postal codes:"]) + if events.get("lookups", 0) and not events.get("upstream_attempts", 0): + lines.extend(["", "ATTENTION: All recorded lookups were served without new upstream attempts. Check the enrichment queue for repeated postal codes."]) + lines.extend(["", "Failed lookup attempts in this reporting window:"]) if failures: for row in failures: lines.append(f"- {row.get('postal_code')}: {row.get('message') or row.get('status')}") diff --git a/cache_manager_new.py b/cache_manager_new.py index 659412b..d6df07e 100644 --- a/cache_manager_new.py +++ b/cache_manager_new.py @@ -16,6 +16,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__) @@ -682,6 +684,12 @@ 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")) + data["average_household_net_worth_amount"] = amount + data["net_worth_source"] = "Historical segment reference" if data.get("average_household_net_worth") else "Unavailable" data.update( { "cached_at": row["cached_at"], @@ -723,7 +731,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', ?) """, @@ -738,7 +754,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) @@ -757,11 +781,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", @@ -774,6 +830,7 @@ 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", "education", "urbanity", "occupation", @@ -796,8 +853,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: diff --git a/docs/2026-09-06-upgrade.md b/docs/2026-09-06-upgrade.md new file mode 100644 index 0000000..285a435 --- /dev/null +++ b/docs/2026-09-06-upgrade.md @@ -0,0 +1,33 @@ +# PRIZM operational upgrade + +## Verified production findings (September 6, 2026) + +Railway production runs main 07816f0. The Salesforce org is Wilderness Committee (00DAm0000008bmtMAA), not a sandbox. August 31–September 6: 10 successful cached responses daily, 70 total, zero upstream attempts. Cache: 161 active successes, 2 expired records, newest cached_at July 3. Salesforce daily queueables completed without platform errors. + +The daily selector treats missing net worth as incomplete enrichment. At least 140 V0N 1T0 households (segment 66) and 70 V0G 1M0 households (segment 50) already have scores, income and property values but no net worth. The unavailable optional field repeatedly selects those postal codes. Partial DML results are also currently ignored. + +Live PRIZM reference rows still include Dwelling Type, but have no net-worth column. All 161 active cache rows have home_type. Net worth is supplemented from a partial legacy segment map; it is not a current postal-code-specific wealth measurement. Certain dwelling strings cannot map exactly to the Salesforce picklist. + +The root/dashboard, Basic authentication and cache CSV export already exist. Weekly email is scheduled in Salesforce, Monday 08:00 America/Los_Angeles, to dena@wildernesscommittee.org. August 31 queueable completed; this is not evidence of inbox delivery. Railway SMTP is not configured and is unnecessary for the existing Salesforce sender. + +## Acceptance criteria + +- Unavailable net worth alone cannot consume future daily batches; never erase an existing value with an unavailable result. Transient upstream failures must remain retryable and partial Salesforce update failures must be visible. +- Dashboard is password protected even when API key is absent. Dashboard credentials grant read-only reporting access, not email-send permission. +- Dashboard shows daily attempts, cache hits, upstream successes/failures, unique codes and an explicit stalled-capture warning. All cached records are reachable by pagination; each row exposes full record details. +- Show failed lookup events even when quota/network errors were deliberately not cached. Export failures and the entire cache without a 5,000-row truncation. Neutralize CSV spreadsheet formulas. +- Dashboard and weekly report describe net-worth values as historical segment references, distinguish new-to-history captures from refreshed upstream successes, and use a clearly labelled rolling seven-day UTC reporting window. +- Preserve existing scheduled weekly email. Validate Python behavior and Salesforce deployment without applying production changes. Request approval for the concrete deployment when ready. + +## Open verification + +The configured API batch maximum is 10, but this is not proof of upstream entitlement. Current public frontend says search-limit reached without documenting a numeric three-per-day cap. Do not infer a usable quota from cache hits, probe around restrictions, or invent missing net-worth data. Confirm licensed allowance with Dena/Environics before changing upstream volume. + +## Validation and rollout + +- 19 Python unit tests pass, including authentication without an API key, uncached failure history, pagination, CSV formula escaping and a 5,001-row export. +- Ruff passes for the Python CI targets. Browser review confirms authenticated dashboard rendering, daily metrics, failures and weekly preview against an isolated fixture database. +- Salesforce validation-only deployment `0AfN3000001kBntKAE` succeeded with six Apex tests and zero errors; no production code or household records were changed. Regression tests cover optional net worth, retaining a known value and retryable failures. +- Roll out the Railway API first (adds the retryable error flag), verify protected dashboard/report/export endpoints, then deploy the validated Salesforce classes. Preserve the existing schedules and credentials. +- Daily DML changes from silently ignored partial results to atomic updates: one household validation failure now rolls back that batch and marks the queueable failed. Inspect the failure before retrying. This intentionally favors visible, consistent outcomes. +- After approval/deployment, observe the next daily run for new postal codes, upstream availability and successful Salesforce updates. This cannot be proven by the validation-only deployment. If upstream quota is unavailable, obtain the licensed allowance; do not bypass the restriction. diff --git a/salesforce/classes/WCPrizmDailyEnrichment.cls b/salesforce/classes/WCPrizmDailyEnrichment.cls index 975c105..82860a6 100644 --- a/salesforce/classes/WCPrizmDailyEnrichment.cls +++ b/salesforce/classes/WCPrizmDailyEnrichment.cls @@ -92,7 +92,6 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data AND ( PRIZM_Score__c = null OR Average_Income_in_Postal_Code__c = null - OR Average_Net_Worth_in_Postal_Code__c = null OR Property_ownership_and_value__c = null ) LIMIT 9000 @@ -107,9 +106,15 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data if (result.status == 'success') { updateAccount.PRIZM_Score__c = parseDecimal(result.segmentNumber); updateAccount.Average_Income_in_Postal_Code__c = mapIncome(result.averageHouseholdIncome); - updateAccount.Average_Net_Worth_in_Postal_Code__c = mapNetWorth(result.averageHouseholdNetWorth, result.averageHouseholdNetWorthAmount); + // Net worth is optional historical reference data, not a completion requirement. + String netWorth = mapNetWorth(result.averageHouseholdNetWorth, result.averageHouseholdNetWorthAmount); + if (netWorth != null) { + updateAccount.Average_Net_Worth_in_Postal_Code__c = netWorth; + } updateAccount.WC_Neighbourhood_description__c = firstNonBlank(new List{ result.whoTheyAre, result.segmentDescription }); updateAccount.Property_ownership_and_value__c = mapHomeType(result.homeType); + } else if (result.retryable) { + continue; } else { updateAccount.PRIZM_Score__c = 0; updateAccount.WC_Neighbourhood_description__c = firstNonBlank(new List{ @@ -122,7 +127,9 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data } if (!updates.isEmpty()) { - Database.update(updates, false); + // Fail the job visibly and roll back the batch if any household update fails. + // Ignoring SaveResult previously made unsuccessful updates look completed. + update updates; } } @@ -140,7 +147,6 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data AND ( PRIZM_Score__c = null OR Average_Income_in_Postal_Code__c = null - OR Average_Net_Worth_in_Postal_Code__c = null OR Property_ownership_and_value__c = null ) GROUP BY BillingPostalCode @@ -352,6 +358,7 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data private String homeType; private String status; private String message; + private Boolean retryable; private PrizmResult(Map row) { postalCode = (String) row.get('postal_code'); @@ -363,6 +370,7 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data averageHouseholdNetWorthAmount = decimalValue(row.get('average_household_net_worth_amount')); homeType = (String) row.get('home_type'); status = (String) row.get('status'); + retryable = row.get('retryable') == true; message = (String) row.get('message'); } diff --git a/salesforce/classes/WCPrizmDailyEnrichmentTest.cls b/salesforce/classes/WCPrizmDailyEnrichmentTest.cls index 8161a4c..818491b 100644 --- a/salesforce/classes/WCPrizmDailyEnrichmentTest.cls +++ b/salesforce/classes/WCPrizmDailyEnrichmentTest.cls @@ -100,4 +100,60 @@ private class WCPrizmDailyEnrichmentTest { System.assertEquals('V8A 2P4', WCPrizmDailyEnrichment.normalizePostalCode('v8a2p4')); } + + @IsTest + static void missingNetWorthDoesNotBlockQueue() { + Id householdType = [SELECT Id FROM RecordType WHERE SobjectType = 'Account' AND DeveloperName = 'HH_Account' LIMIT 1].Id; + Account enriched = new Account( + Name = 'Optional net worth unavailable', RecordTypeId = householdType, + BillingStreet = '1 Test St', BillingCity = 'Powell River', BillingCountryCode = 'CA', + BillingStateCode = 'BC', BillingPostalCode = 'V0N 1T0', PRIZM_Score__c = 66, + Average_Income_in_Postal_Code__c = '$0 to $110,593', + Property_ownership_and_value__c = 'None or Unverifiable' + ); + insert enriched; + System.assertEquals(0, WCPrizmDailyEnrichment.selectPostalCodes().size(), + 'An assigned household must not be selected solely because optional net worth is unavailable.'); + Account pending = new Account( + Name = 'New postal code', RecordTypeId = householdType, + BillingStreet = '2 Test St', BillingCity = 'Powell River', BillingCountryCode = 'CA', + BillingStateCode = 'BC', BillingPostalCode = 'V8A 2P4' + ); + insert pending; + System.assertEquals(1, WCPrizmDailyEnrichment.selectPostalCodes().size()); + } + + private class MissingNetWorthMock implements HttpCalloutMock { + public HTTPResponse respond(HTTPRequest request) { + HttpResponse response = new HttpResponse(); + response.setStatusCode(200); + response.setBody('{"results":[{"postal_code":"V8A 2P4","segment_number":"66","average_household_income":"$98,721","home_type":"Single Detached","status":"success"},{"postal_code":"V5U 1J4","status":"error","retryable":true,"message":"quota unavailable"}]}'); + return response; + } + } + + @IsTest + static void preservesKnownNetWorthAndRetriesTransientFailures() { + Id householdType = [SELECT Id FROM RecordType WHERE SobjectType = 'Account' AND DeveloperName = 'HH_Account' LIMIT 1].Id; + Account known = new Account( + Name = 'Preserve historical value', RecordTypeId = householdType, + BillingStreet = '1 Test St', BillingCity = 'Powell River', BillingCountryCode = 'CA', + BillingStateCode = 'BC', BillingPostalCode = 'V8A 2P4', + Average_Net_Worth_in_Postal_Code__c = '$600K to $1M' + ); + Account retry = new Account( + Name = 'Retry quota failure', RecordTypeId = householdType, + BillingStreet = '2 Test St', BillingCity = 'Vancouver', BillingCountryCode = 'CA', + BillingStateCode = 'BC', BillingPostalCode = 'V5U 1J4' + ); + insert new List{ known, retry }; + Test.setMock(HttpCalloutMock.class, new MissingNetWorthMock()); + Test.startTest(); + WCPrizmDailyEnrichment.runDaily(); + Test.stopTest(); + known = [SELECT Average_Net_Worth_in_Postal_Code__c FROM Account WHERE Id = :known.Id]; + retry = [SELECT PRIZM_Score__c FROM Account WHERE Id = :retry.Id]; + System.assertEquals('$600K to $1M', known.Average_Net_Worth_in_Postal_Code__c); + System.assertEquals(null, retry.PRIZM_Score__c, 'Temporary quota failures must not permanently exclude households.'); + } } \ No newline at end of file diff --git a/salesforce/classes/WCPrizmWeeklyReportEmail.cls b/salesforce/classes/WCPrizmWeeklyReportEmail.cls index afe7cc8..d3adced 100644 --- a/salesforce/classes/WCPrizmWeeklyReportEmail.cls +++ b/salesforce/classes/WCPrizmWeeklyReportEmail.cls @@ -37,7 +37,10 @@ public without sharing class WCPrizmWeeklyReportEmail implements Schedulable { email.setToAddresses(new List{ DEFAULT_RECIPIENT }); email.setSubject(subject); email.setPlainTextBody(body); - Messaging.sendEmail(new List{ email }); + Messaging.SendEmailResult[] sent = Messaging.sendEmail(new List{ email }); + if (!sent[0].isSuccess()) { + throw new CalloutException('PRIZM report email failed: ' + sent[0].getErrors()[0].getMessage()); + } } private String stringValue(Object value) { diff --git a/test_prizm_api.py b/test_prizm_api.py index 30880cd..10f4421 100644 --- a/test_prizm_api.py +++ b/test_prizm_api.py @@ -136,5 +136,106 @@ def test_api_key_protects_all_routes_except_health(self): os.environ.pop("PRIZM_API_KEY", None) + + +class TestOperationalReporting(unittest.TestCase): + def setUp(self): + import tempfile + from cache_manager_new import CacheManager + self.directory = tempfile.TemporaryDirectory() + self.addCleanup(self.directory.cleanup) + self.cache = CacheManager(os.path.join(self.directory.name, 'cache.db')) + self.override = patch('app.cache_manager', self.cache) + self.override.start() + self.addCleanup(self.override.stop) + self.env = patch.dict(os.environ, {'PRIZM_API_KEY': 'api-secret', 'DASHBOARD_PASSWORD': 'dashboard-secret', 'DASHBOARD_USERNAME': 'dena'}, clear=True) + self.env.start() + self.addCleanup(self.env.stop) + import base64 + self.auth = {'Authorization': 'Basic ' + base64.b64encode(b'dena:dashboard-secret').decode()} + self.client = app.test_client() + + def test_dashboard_fails_closed_without_either_secret(self): + with patch.dict(os.environ, {}, clear=True): + for path in ['/', '/dashboard', '/api/dashboard/summary', '/api/cache/entries', '/api/cache/export.csv', '/api/lookups', '/api/lookups/failures.csv', '/api/reports/weekly']: + self.assertEqual(self.client.get(path).status_code, 401, path) + self.assertEqual(self.client.get('/health').status_code, 200) + + def test_dashboard_password_required_when_api_key_absent(self): + os.environ.pop('PRIZM_API_KEY') + self.assertEqual(self.client.get('/').status_code, 401) + self.assertEqual(self.client.get('/', headers=self.auth).status_code, 200) + + def test_dashboard_credentials_cannot_send_email_or_modify_cache(self): + self.assertEqual(self.client.get('/api/reports/weekly', headers=self.auth).status_code, 200) + self.assertEqual(self.client.post('/api/reports/weekly/send', headers=self.auth).status_code, 401) + self.assertEqual(self.client.post('/api/cache/clear', headers=self.auth).status_code, 401) + + def test_events_distinguish_repeated_cache_hits_and_fresh_successes(self): + self.cache.record_lookup_event('V8A 0A8', 'success', 'cache', from_cache=True) + self.cache.record_lookup_event('V8A 0A8', 'success', 'cache', from_cache=True) + self.cache.record_lookup_event('V8A 0A8', 'success', 'upstream') + self.cache.record_lookup_event('V8A 2P4', 'success', 'upstream') + self.cache.record_lookup_event('M5V 3L9', 'error', 'upstream', message='quota unavailable') + result = self.cache.get_lookup_event_summary() + self.assertEqual(result['lookups'], 5) + self.assertEqual(result['unique_postal_codes'], 3) + self.assertEqual(result['cache_hits'], 2) + self.assertEqual(result['upstream_successful'], 2) + self.assertEqual(result['upstream_failed'], 1) + self.assertEqual(result['newly_captured'], 1) + self.assertEqual(result['by_day'][0]['newly_captured'], 1) + summary = self.client.get('/api/dashboard/summary', headers=self.auth).get_json() + self.assertEqual(summary['recent_failures'][0]['message'], 'quota unavailable') + self.assertEqual(self.cache.list_cache_entries(), []) + + def test_failure_history_paginates_and_csv_neutralizes_formulas(self): + import csv + import io + self.cache.record_lookup_event('V8A 0A8', 'error', 'upstream', message='=HYPERLINK("bad")') + self.cache.record_lookup_event('V8A 2P4', 'invalid', 'validation', message='invalid') + first = self.client.get('/api/lookups?failures=1&limit=1', headers=self.auth).get_json()['entries'] + second = self.client.get('/api/lookups?failures=1&limit=1&offset=1', headers=self.auth).get_json()['entries'] + self.assertNotEqual(first[0]['id'], second[0]['id']) + rows = list(csv.DictReader(io.StringIO(self.cache.export_failures_csv()))) + self.assertEqual(len(rows), 2) + self.assertTrue(rows[1]['message'].startswith("'=")) + + def test_csv_exports_more_than_5000_entries_and_marks_historical_values(self): + import csv + import io + self.cache.cache_data('V8A 0A8', LOOKUP_RESULT, custom_duration_days=10) + with self.cache._connect() as conn: + columns = [r[1] for r in conn.execute('PRAGMA table_info(postal_code_cache)') if r[1] not in ('postal_code', 'id')] + names = ','.join(columns) + conn.executemany(f'INSERT INTO postal_code_cache (postal_code,{names}) SELECT ?,{names} FROM postal_code_cache WHERE postal_code = ?', [(f'X{i:05}', 'V8A 0A8') for i in range(5000)]) + conn.commit() + rows = list(csv.DictReader(io.StringIO(self.cache.export_cache_csv()))) + self.assertEqual(len(rows), 5001) + self.assertEqual(rows[0]['net_worth_source'], 'Historical segment reference') + + def test_report_excludes_old_failures_and_warns_about_cache_only_activity(self): + from app import build_weekly_report + self.cache.record_lookup_event('V8A 0A8', 'error', 'upstream', message='old failure') + with self.cache._connect() as conn: + conn.execute("UPDATE lookup_events SET requested_at = datetime('now', '-10 days')") + conn.commit() + self.cache.record_lookup_event('V8A 2P4', 'success', 'cache', from_cache=True) + report = build_weekly_report() + self.assertIn('ATTENTION', report['body']) + self.assertNotIn('old failure', report['body']) + self.assertIn('Unique postal codes requested: 1', report['body']) + report = build_weekly_report(days=14) + self.assertIn('old failure', report['body']) + self.assertIn('Lookups recorded: 2', report['body']) + + @patch('app.prizm_client.lookup', side_effect=PrizmLookupError('quota unavailable')) + def test_transient_failure_is_retryable_and_visible_without_cache(self, lookup): + response = self.client.get('/api/prizm?postal_code=V8A0A8', headers={'X-API-Key': 'api-secret'}).get_json() + self.assertTrue(response['retryable']) + self.assertEqual(self.cache.list_cache_entries(), []) + self.assertEqual(len(self.cache.list_lookup_events(failures_only=True)), 1) + + if __name__ == "__main__": unittest.main() From d5bb1f0d22ce235ea615933a5bb99ffbbb8287b8 Mon Sep 17 00:00:00 2001 From: heybeaux Date: Sun, 6 Sep 2026 16:54:02 -0700 Subject: [PATCH 2/5] fix: collect new major-donor postal codes from a durable cohort queue --- .github/workflows/ci.yml | 4 +- app.py | 70 ++++++++- cache_manager_new.py | 21 +++ cohort_queue.py | 148 ++++++++++++++++++ docs/2026-09-06-upgrade.md | 77 +++++++-- salesforce/classes/WCPrizmDailyEnrichment.cls | 112 +++---------- .../classes/WCPrizmDailyEnrichmentTest.cls | 41 +++-- test_cohort_queue.py | 98 ++++++++++++ 8 files changed, 440 insertions(+), 131 deletions(-) create mode 100644 cohort_queue.py create mode 100644 test_cohort_queue.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d124075..04ab3b9 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -30,10 +30,10 @@ jobs: pip install -r requirements.txt ruff - 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 test_prizm_api.py test_cohort_queue.py - name: Test - run: python -m unittest test_prizm_api.py + run: python -m unittest test_prizm_api.py test_cohort_queue.py node: name: Node tests and build diff --git a/app.py b/app.py index 11fc73d..21e711c 100644 --- a/app.py +++ b/app.py @@ -13,6 +13,7 @@ from flask import Flask, Response, jsonify, make_response, request from cache_manager_new import cache_manager +from cohort_queue import CohortQueue from prizm_client import PrizmClient, PrizmLookupError, normalize_postal_code from segment_net_worth import average_household_net_worth, average_household_net_worth_amount @@ -87,6 +88,9 @@ +

Major-donor coverage

+

Loading…

+ Export major-donor coverage CSV

Daily lookup outcomes

Trailing 7 days, grouped by UTC date. Cache hits do not represent new captures.
@@ -133,6 +137,8 @@ text('failed', fmt.format((breakdown.error || 0) + (breakdown.invalid || 0))); text('weekLookups', fmt.format(week.lookups || 0)); text('captureWarning', week.lookups && !week.upstream_attempts ? 'No new upstream lookups in the last 7 days. Check whether the enrichment queue is repeating already captured postal codes.' : ''); + const cohort = data.cohort || {}; + text('cohortProgress', cohort.configured ? `${cohort.captured} / ${cohort.unique_codes} codes captured · ${cohort.pending} pending · ${cohort.failed} failed · ${cohort.incomplete} incomplete. ${cohort.source_rows} source accounts; ${cohort.excluded_rows} address exceptions. Daily target: 10 new codes; no automatic repeats.` : 'Cohort has not been imported.'); const counts = week.by_day || []; document.getElementById('dailyCounts').innerHTML = counts.length ? counts.map(row => `
${esc(row.day)}: ${fmt.format(row.lookups || 0)} lookups · ${fmt.format(row.cache_hits || 0)} cache hits · ${fmt.format(row.upstream_attempts || 0)} upstream attempts · ${fmt.format(row.upstream_successful || 0)} upstream success · ${fmt.format(row.failed || 0)} failed · ${fmt.format(row.unique_postal_codes || 0)} unique codes · ${fmt.format(row.newly_captured || 0)} first captures in recorded history
` @@ -254,7 +260,7 @@ def require_authentication(): if request.path == "/health" or request.method == "OPTIONS": return None - dashboard_paths = {"/", "/dashboard", "/api/dashboard/summary", "/api/cache/entries", "/api/cache/export.csv", "/api/reports/weekly", "/api/lookups", "/api/lookups/failures.csv"} + dashboard_paths = {"/", "/dashboard", "/api/dashboard/summary", "/api/cache/entries", "/api/cache/export.csv", "/api/reports/weekly", "/api/lookups", "/api/lookups/failures.csv", "/api/cohort", "/api/cohort/export.csv"} if request.path in dashboard_paths: if has_valid_dashboard_auth() or (os.environ.get("PRIZM_API_KEY") and has_valid_api_key()): return None @@ -415,6 +421,59 @@ def get_batch_prizm(): ) +def cohort_queue(): + return CohortQueue(cache_manager.db_path) + + +@app.route("/api/cohort", methods=["GET"]) +def cohort_status(): + queue = cohort_queue() + return jsonify({"summary": queue.summary(), "entries": queue.entries()}) + + +@app.route("/api/cohort/export.csv", methods=["GET"]) +def cohort_export(): + response = Response(cohort_queue().export_csv(), mimetype="text/csv") + response.headers["Content-Disposition"] = "attachment; filename=major-donor-coverage.csv" + return response + + +@app.route("/api/cohort/import", methods=["POST"]) +def import_cohort(): + if not os.environ.get("PRIZM_API_KEY") or not has_valid_api_key(): + return jsonify({"error": "Import requires an API key"}), 401 + if request.content_length is None or request.content_length > 1_000_000: + return jsonify({"error": "A CSV body under 1 MB is required"}), 400 + try: + import io + summary = cohort_queue().import_accounts(io.StringIO(request.get_data().decode("utf-8-sig"))) + except (ValueError, UnicodeError): + return jsonify({"error": "Invalid Canadian billing-address CSV; cohort unchanged"}), 400 + return jsonify(summary) + + +@app.route("/api/cohort/capture", methods=["POST"]) +def capture_cohort(): + # Collection is a mutation and always requires an explicit API key. + if not os.environ.get("PRIZM_API_KEY") or not has_valid_api_key(): + return jsonify({"error": "Collection requires an API key"}), 401 + queue = cohort_queue() + if not queue.summary()["configured"]: + return jsonify({"error": "Major-donor cohort has not been imported"}), 503 + results = [] + # One code per call keeps Salesforce callouts below transaction time limits. + # Its queueable chains ten calls; the database also enforces the daily ceiling. + for _ in range(1): + code = queue.claim() + if code is None: + break + # Claim persists before network I/O: retries/concurrent requests cannot repeat a code. + result = get_prizm_code(code, endpoint="cohort") + queue.finish(code, result) + results.append(result) + return jsonify({"results": results, "summary": queue.summary()}) + + @app.route("/api/segments", methods=["GET"]) def get_segments(): return jsonify({"status": "success", "segments": prizm_client.get_all_segments()}) @@ -422,7 +481,9 @@ def get_segments(): @app.route("/api/dashboard/summary", methods=["GET"]) def dashboard_summary(): - return jsonify(cache_manager.get_dashboard_summary()) + summary = cache_manager.get_dashboard_summary() + summary["cohort"] = cohort_queue().summary() + return jsonify(summary) @app.route("/api/lookups", methods=["GET"]) @@ -535,6 +596,7 @@ def build_weekly_report(days: int = 7) -> Dict[str, Any]: summary = cache_manager.get_dashboard_summary() events = cache_manager.get_lookup_event_summary(days) summary["lookup_events_7d"] = events + summary["cohort"] = cohort_queue().summary() stats = summary.get("cache_stats", {}) daily_counts = events.get("by_day", []) failures = cache_manager.list_lookup_events(failures_only=True, days=days, limit=20) @@ -562,6 +624,10 @@ def build_weekly_report(days: int = 7) -> Dict[str, Any]: "", "Daily lookup outcomes (UTC):", ] + cohort = summary["cohort"] + lines.extend(["", "Major-donor Canadian postal-code coverage:", + f"Captured: {cohort['captured']}; pending: {cohort['pending']}; failed: {cohort['failed']}; incomplete: {cohort['incomplete']}", + "Daily target: 10 never-attempted codes. Failed/incomplete attempts require review; no automatic repeats.", ""]) if daily_counts: for row in daily_counts: lines.append(f"- {row.get('day')}: {row.get('lookups', 0)} lookups, {row.get('cache_hits', 0)} cached, {row.get('upstream_successful', 0)} upstream successes, {row.get('failed', 0)} failures") diff --git a/cache_manager_new.py b/cache_manager_new.py index d6df07e..7b305e9 100644 --- a/cache_manager_new.py +++ b/cache_manager_new.py @@ -283,6 +283,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) diff --git a/cohort_queue.py b/cohort_queue.py new file mode 100644 index 0000000..fd1aa95 --- /dev/null +++ b/cohort_queue.py @@ -0,0 +1,148 @@ +"""Private cohort collection, independent of Salesforce field completeness and cache TTL.""" +import csv +import io +import json +import re +import sqlite3 +from contextlib import contextmanager +from datetime import datetime +from zoneinfo import ZoneInfo + + +def postal_code(value): + value = re.sub(r'\s+', '', (value or '').upper()) + if re.fullmatch(r'[ABCEGHJ-NPRSTVXY][0-9][ABCEGHJ-NPRSTV-Z][0-9][ABCEGHJ-NPRSTV-Z][0-9]', value): + return value[:3] + ' ' + value[3:] + return None + + +class CohortQueue: + def __init__(self, db_path): + self.db_path = db_path + with self.connect() as db: + db.executescript(''' + CREATE TABLE IF NOT EXISTS capture_history ( + postal_code TEXT PRIMARY KEY, captured_at TEXT NOT NULL + ); + CREATE TABLE IF NOT EXISTS donor_cohort ( + postal_code TEXT PRIMARY KEY, account_rows INTEGER NOT NULL + ); + CREATE TABLE IF NOT EXISTS cohort_attempts ( + postal_code TEXT PRIMARY KEY, day TEXT NOT NULL, + started_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, + result_json TEXT + ); + CREATE TABLE IF NOT EXISTS cohort_metadata ( + id INTEGER PRIMARY KEY CHECK(id=1), source_rows INTEGER NOT NULL, + excluded_rows INTEGER NOT NULL + ); + ''') + # Seed all historical successes, including expired cache and retained events. + for query in ( + "SELECT postal_code, cached_at FROM postal_code_cache WHERE status='success'", + "SELECT postal_code, MIN(requested_at) FROM lookup_events WHERE status='success' GROUP BY postal_code", + ): + for raw, captured in db.execute(query).fetchall(): + code = postal_code(raw) + if code: + db.execute('INSERT OR IGNORE INTO capture_history VALUES (?, ?)', (code, str(captured))) + + @contextmanager + def connect(self): + db = sqlite3.connect(self.db_path, timeout=30) + db.row_factory = sqlite3.Row + try: + with db: + yield db + finally: + db.close() + + def import_accounts(self, stream): + reader = csv.DictReader(stream) + required = {'Billing Country Code', 'Billing Zip/Postal Code'} + if not required.issubset(reader.fieldnames or []): + raise ValueError('Salesforce billing country and postal-code columns are required') + counts, total, excluded = {}, 0, 0 + for row in reader: + total += 1 + code = postal_code(row['Billing Zip/Postal Code']) + if (row.get('Billing Country Code') or '').strip().upper() != 'CA' or not code: + excluded += 1 + continue + counts[code] = counts.get(code, 0) + 1 + if not counts: + raise ValueError('No eligible Canadian billing postal codes; existing cohort preserved') + with self.connect() as db: + db.execute('BEGIN IMMEDIATE') + # Exact membership replacement; retain capture and attempt history permanently. + db.execute('DELETE FROM donor_cohort') + db.executemany('INSERT INTO donor_cohort VALUES (?, ?)', counts.items()) + db.execute('INSERT OR REPLACE INTO cohort_metadata VALUES (1, ?, ?)', (total, excluded)) + return self.summary() + + def claim(self, day=None): + day = day or datetime.now(ZoneInfo('America/Vancouver')).date().isoformat() + with self.connect() as db: + db.execute('BEGIN IMMEDIATE') + if db.execute('SELECT COUNT(*) FROM cohort_attempts WHERE day=?', (day,)).fetchone()[0] >= 10: + return None + row = db.execute(''' + SELECT c.postal_code FROM donor_cohort c + LEFT JOIN capture_history h USING(postal_code) + LEFT JOIN cohort_attempts a USING(postal_code) + WHERE h.postal_code IS NULL AND a.postal_code IS NULL + ORDER BY c.account_rows DESC, c.postal_code LIMIT 1 + ''').fetchone() + if row is None: + return None + db.execute('INSERT INTO cohort_attempts (postal_code, day) VALUES (?, ?)', (row[0], day)) + return row[0] + + def finish(self, code, result): + with self.connect() as db: + db.execute('UPDATE cohort_attempts SET result_json=? WHERE postal_code=?', (json.dumps(result), code)) + if result.get('status') == 'success': + db.execute("INSERT OR IGNORE INTO capture_history VALUES (?, CURRENT_TIMESTAMP)", (code,)) + + def entries(self): + with self.connect() as db: + return [dict(row) for row in db.execute(''' + SELECT c.postal_code, c.account_rows, h.captured_at, a.started_at, + CASE WHEN h.postal_code IS NOT NULL THEN 'captured' + WHEN a.result_json IS NOT NULL THEN 'failed' + WHEN a.postal_code IS NOT NULL THEN 'incomplete' + ELSE 'pending' END AS coverage, + a.result_json + FROM donor_cohort c LEFT JOIN capture_history h USING(postal_code) + LEFT JOIN cohort_attempts a USING(postal_code) + ORDER BY c.account_rows DESC, c.postal_code + ''')] + + def summary(self): + rows = self.entries() + with self.connect() as db: + meta = db.execute('SELECT source_rows, excluded_rows FROM cohort_metadata').fetchone() + return dict(meta or {}, configured=meta is not None, unique_codes=len(rows), + eligible_accounts=sum(row['account_rows'] for row in rows), + **{key: sum(row['coverage'] == key for row in rows) + for key in ('captured', 'pending', 'failed', 'incomplete')}) + + def export_csv(self): + stream = io.StringIO() + writer = csv.DictWriter(stream, fieldnames=['postal_code', 'account_rows', 'coverage', 'captured_at', 'started_at']) + writer.writeheader() + for row in self.entries(): + writer.writerow({key: row[key] for key in writer.fieldnames}) + return stream.getvalue() + + +if __name__ == '__main__': + import argparse + from cache_manager_new import CacheManager + parser = argparse.ArgumentParser(description='Import a private Salesforce account export; never scrape or update Salesforce.') + parser.add_argument('csv_file') + parser.add_argument('--db', required=True) + args = parser.parse_args() + CacheManager(args.db) + with open(args.csv_file, encoding='utf-8-sig', newline='') as stream: + print(json.dumps(CohortQueue(args.db).import_accounts(stream), indent=2)) diff --git a/docs/2026-09-06-upgrade.md b/docs/2026-09-06-upgrade.md index 285a435..833558a 100644 --- a/docs/2026-09-06-upgrade.md +++ b/docs/2026-09-06-upgrade.md @@ -10,24 +10,69 @@ Live PRIZM reference rows still include Dwelling Type, but have no net-worth col The root/dashboard, Basic authentication and cache CSV export already exist. Weekly email is scheduled in Salesforce, Monday 08:00 America/Los_Angeles, to dena@wildernesscommittee.org. August 31 queueable completed; this is not evidence of inbox delivery. Railway SMTP is not configured and is unnecessary for the existing Salesforce sender. -## Acceptance criteria +## Final collection policy confirmed by Beaux -- Unavailable net worth alone cannot consume future daily batches; never erase an existing value with an unavailable result. Transient upstream failures must remain retryable and partial Salesforce update failures must be visible. -- Dashboard is password protected even when API key is absent. Dashboard credentials grant read-only reporting access, not email-send permission. -- Dashboard shows daily attempts, cache hits, upstream successes/failures, unique codes and an explicit stalled-capture warning. All cached records are reachable by pagination; each row exposes full record details. -- Show failed lookup events even when quota/network errors were deliberately not cached. Export failures and the entire cache without a 5,000-row truncation. Neutralize CSV spreadsheet formulas. -- Dashboard and weekly report describe net-worth values as historical segment references, distinguish new-to-history captures from refreshed upstream successes, and use a clearly labelled rolling seven-day UTC reporting window. -- Preserve existing scheduled weekly email. Validate Python behavior and Salesforce deployment without applying production changes. Request approval for the concrete deployment when ready. +Use the supplied 513 Account records as the complete major-donor cohort. Collect +Canadian billing postal codes only. Daily target is ten **new** codes, not three. +Never automatically select an already enriched code again, including after cache +expiry. The supplied export contains 503 eligible accounts, 484 unique codes; +106 codes already captured cover 112 accounts, leaving 378 codes across 391 +accounts. Nine non-Canadian addresses and one missing Canadian billing postal +code remain address exceptions. Shipping addresses are not used. -## Open verification +## Implemented locally -The configured API batch maximum is 10, but this is not proof of upstream entitlement. Current public frontend says search-limit reached without documenting a numeric three-per-day cap. Do not infer a usable quota from cache hits, probe around restrictions, or invent missing net-worth data. Confirm licensed allowance with Dena/Environics before changing upstream volume. +- Railway owns a private cohort table and a permanent capture-history table seeded + from existing successes, including expired cache and retained lookup events. + Database triggers retain future captures even when cache rows are deleted. +- Atomic claims select never-attempted, never-captured cohort codes, ordered by + account frequency then postal code. At most ten claims per Vancouver calendar + day across repeated requests, concurrent workers and service restarts. +- Each API call claims one code. The existing Salesforce scheduler starts a chain + of up to ten queueable jobs, each with one callout and matching household updates. + This avoids accumulating ten upstream lookups inside one callout transaction. + The general Salesforce candidate selector is removed. An empty or unconfigured + cohort never falls back to general household selection. +- Failed and interrupted attempts require review and are not automatically retried. + Failed captures do not count as enriched. Ten is an attempt target; upstream + restrictions/errors may reduce successful captures. This configuration is not + proof of the provider's current licensed allowance. +- Protected dashboard and CSV export show cohort capture/pending/failure/incomplete + counts; the existing Salesforce weekly email includes cohort coverage. Import + and capture mutations require an API key, never dashboard credentials alone. +- Capture progress does not prove Salesforce record updates. Existing household + mapping is retained; organizations are captured but not updated by that household + mapper. Missing historical net worth cannot drive collection or erase a known value. -## Validation and rollout +## Local validation -- 19 Python unit tests pass, including authentication without an API key, uncached failure history, pagination, CSV formula escaping and a 5,001-row export. -- Ruff passes for the Python CI targets. Browser review confirms authenticated dashboard rendering, daily metrics, failures and weekly preview against an isolated fixture database. -- Salesforce validation-only deployment `0AfN3000001kBntKAE` succeeded with six Apex tests and zero errors; no production code or household records were changed. Regression tests cover optional net worth, retaining a known value and retryable failures. -- Roll out the Railway API first (adds the retryable error flag), verify protected dashboard/report/export endpoints, then deploy the validated Salesforce classes. Preserve the existing schedules and credentials. -- Daily DML changes from silently ignored partial results to atomic updates: one household validation failure now rolls back that batch and marks the queueable failed. Inspect the failure before retrying. This intentionally favors visible, consistent outcomes. -- After approval/deployment, observe the next daily run for new postal codes, upstream availability and successful Salesforce updates. This cannot be proven by the validation-only deployment. If upstream quota is unavailable, obtain the licensed allowance; do not bypass the restriction. +Python regression and concurrency tests cover daily limits, duplicate prevention, +restart behavior, expiry/deletion, account normalization/exclusions, imports and +API authentication. Browser verified desktop/mobile cohort progress (106/484), +weekly preview and CSV (484 data rows) against an isolated production-cache copy. +26 Python tests and Ruff pass. Salesforce validation-only deployment +`0AfN3000001kBr7KAE` succeeded with six tests and zero component/test errors. +The first validation exposed recursive mock responses in the schedulable test; +the scheduler test now returns an exhausted cohort. The final validation client +timed out, but the server-side deployment report confirms success. + +## Deployment and operational limits + +No production change has been made. After approval, deploy Railway, import the +private two-column country/postal-code payload through `/api/cohort/import`, verify +106 captured / 378 pending, then deploy validated Salesforce classes. Preserve +existing schedule and credentials. Import/history schema are additive; retain the +original source and a database backup before rollout. No donor addresses or cohort +membership should enter GitHub. + +Observe the first scheduled run for ten distinct new attempts and actual captures. +If a Salesforce validation rule rejects an update, that queueable fails visibly; +its API capture remains saved and will not be scraped again. The remainder of that +chain stops. A response/network interruption can similarly leave a captured code +without a Salesforce update. Reconcile such records from stored results without +repeating upstream collection. Incomplete claims must be reviewed before any retry. + +The confirmed recent stall is 70 cache hits and zero upstream attempts August 31– +September 6; latest cache timestamp is July 3. Full historical causality for every +day since July 3 has not been independently reconstructed. Database SSH was +unavailable with the configured local key; no new key was registered. diff --git a/salesforce/classes/WCPrizmDailyEnrichment.cls b/salesforce/classes/WCPrizmDailyEnrichment.cls index 82860a6..501edc2 100644 --- a/salesforce/classes/WCPrizmDailyEnrichment.cls +++ b/salesforce/classes/WCPrizmDailyEnrichment.cls @@ -1,6 +1,4 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Database.AllowsCallouts { - private static final Integer MAX_POSTAL_CODES_PER_RUN = 10; - private static final Integer AGGREGATE_SCAN_LIMIT = 500; private static final Set INCOME_VALUES = new Set{ '$0 to $110,593', @@ -38,7 +36,7 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data @InvocableMethod( label='Run PRIZM Daily Enrichment' - description='Submits the next 10 highest unpopulated Canadian Household Account postal codes to the PRIZM API and updates matching Accounts.' + description='Captures up to 10 new Canadian major-donor postal codes daily and updates matching Household Accounts.' ) public static void runFromFlow(List ignoredInputs) { runQueued(); @@ -49,25 +47,19 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data } public static Id runQueued() { - return System.enqueueJob(new PrizmDailyQueueable()); + return System.enqueueJob(new PrizmDailyQueueable(10)); } - public static void runDaily() { - List buckets = selectPostalCodes(); - if (buckets.isEmpty()) { - return; - } - - List postalCodes = new List(); + public static Boolean runDaily() { + // Railway owns cohort membership and permanent capture history. + // Never derive new capture eligibility from optional Salesforce fields. + Map resultsByPostalCode = fetchPrizmResults(); Map> rawCodesByNormalized = new Map>(); - for (PostalCodeBucket bucket : buckets) { - postalCodes.add(bucket.normalizedPostalCode); - rawCodesByNormalized.put(bucket.normalizedPostalCode, bucket.rawPostalCodes); + for (String code : resultsByPostalCode.keySet()) { + rawCodesByNormalized.put(code, new Set{code, code.replace(' ', '')}); } - - Map resultsByPostalCode = fetchPrizmResults(postalCodes); if (resultsByPostalCode.isEmpty()) { - return; + return false; } Set rawCodesToUpdate = new Set(); @@ -78,7 +70,7 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data } if (rawCodesToUpdate.isEmpty()) { - return; + return false; } List updates = new List(); @@ -131,65 +123,18 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data // Ignoring SaveResult previously made unsuccessful updates look completed. update updates; } + return true; } @TestVisible - private static List selectPostalCodes() { - Map bucketsByPostalCode = new Map(); - - for (AggregateResult row : [ - SELECT BillingPostalCode postalCode, COUNT(Id) accountCount - FROM Account - WHERE RecordType.DeveloperName = 'HH_Account' - AND BillingCountryCode = 'CA' - AND BillingPostalCode != null - AND (PRIZM_Score__c = null OR PRIZM_Score__c != 0) - AND ( - PRIZM_Score__c = null - OR Average_Income_in_Postal_Code__c = null - OR Property_ownership_and_value__c = null - ) - GROUP BY BillingPostalCode - ORDER BY COUNT(Id) DESC - LIMIT :AGGREGATE_SCAN_LIMIT - ]) { - String rawPostalCode = (String) row.get('postalCode'); - String normalizedPostalCode = normalizePostalCode(rawPostalCode); - if (String.isBlank(normalizedPostalCode)) { - continue; - } - - PostalCodeBucket bucket = bucketsByPostalCode.get(normalizedPostalCode); - if (bucket == null) { - bucket = new PostalCodeBucket(normalizedPostalCode); - bucketsByPostalCode.put(normalizedPostalCode, bucket); - } - bucket.count += ((Decimal) row.get('accountCount')).intValue(); - bucket.rawPostalCodes.add(rawPostalCode); - } - - List buckets = bucketsByPostalCode.values(); - buckets.sort(); - - List selected = new List(); - for (PostalCodeBucket bucket : buckets) { - selected.add(bucket); - if (selected.size() == MAX_POSTAL_CODES_PER_RUN) { - break; - } - } - return selected; - } - - @TestVisible - private static Map fetchPrizmResults(List postalCodes) { + private static Map fetchPrizmResults() { HttpRequest request = new HttpRequest(); request.setMethod('POST'); - request.setEndpoint(System.Label.WC_PRIZM_API_Base_URL + '/api/prizm/batch'); + request.setEndpoint(System.Label.WC_PRIZM_API_Base_URL + '/api/cohort/capture'); request.setHeader('Content-Type', 'application/json'); request.setHeader('X-API-Key', System.Label.WC_PRIZM_API_Key); - request.setTimeout(60000); - request.setBody(JSON.serialize(new Map{ 'postal_codes' => postalCodes })); + request.setTimeout(120000); + request.setBody('{}'); HttpResponse response = new Http().send(request); if (response.getStatusCode() < 200 || response.getStatusCode() >= 300) { @@ -322,28 +267,15 @@ public without sharing class WCPrizmDailyEnrichment implements Schedulable, Data return null; } - @TestVisible - private class PostalCodeBucket implements Comparable { - private String normalizedPostalCode; - private Integer count = 0; - private Set rawPostalCodes = new Set(); - - private PostalCodeBucket(String normalizedPostalCode) { - this.normalizedPostalCode = normalizedPostalCode; - } - - public Integer compareTo(Object otherObject) { - PostalCodeBucket other = (PostalCodeBucket) otherObject; - if (count == other.count) { - return normalizedPostalCode.compareTo(other.normalizedPostalCode); - } - return other.count - count; - } - } - private class PrizmDailyQueueable implements Queueable, Database.AllowsCallouts { + private Integer remaining; + private PrizmDailyQueueable(Integer remaining) { + this.remaining = remaining; + } public void execute(QueueableContext context) { - runDaily(); + if (runDaily() && remaining > 1) { + System.enqueueJob(new PrizmDailyQueueable(remaining - 1)); + } } } diff --git a/salesforce/classes/WCPrizmDailyEnrichmentTest.cls b/salesforce/classes/WCPrizmDailyEnrichmentTest.cls index 818491b..a882361 100644 --- a/salesforce/classes/WCPrizmDailyEnrichmentTest.cls +++ b/salesforce/classes/WCPrizmDailyEnrichmentTest.cls @@ -3,7 +3,7 @@ private class WCPrizmDailyEnrichmentTest { private class PrizmMock implements HttpCalloutMock { public HTTPResponse respond(HTTPRequest request) { System.assertEquals('POST', request.getMethod()); - System.assert(request.getEndpoint().endsWith('/api/prizm/batch')); + System.assert(request.getEndpoint().endsWith('/api/cohort/capture')); HttpResponse response = new HttpResponse(); response.setStatusCode(200); @@ -92,7 +92,7 @@ private class WCPrizmDailyEnrichmentTest { @IsTest static void schedulableEntryPointRuns() { - Test.setMock(HttpCalloutMock.class, new PrizmMock()); + Test.setMock(HttpCalloutMock.class, new EmptyCohortMock()); Test.startTest(); System.schedule('WC PRIZM Test Job', '0 0 2 * * ?', new WCPrizmDailyEnrichment()); @@ -101,26 +101,25 @@ private class WCPrizmDailyEnrichmentTest { System.assertEquals('V8A 2P4', WCPrizmDailyEnrichment.normalizePostalCode('v8a2p4')); } + private class EmptyCohortMock implements HttpCalloutMock { + public HTTPResponse respond(HTTPRequest request) { + System.assert(request.getEndpoint().endsWith('/api/cohort/capture')); + System.assertEquals('{}', request.getBody()); + HttpResponse response = new HttpResponse(); + response.setStatusCode(200); + response.setBody('{"results":[],"summary":{"pending":0}}'); + return response; + } + } + @IsTest - static void missingNetWorthDoesNotBlockQueue() { - Id householdType = [SELECT Id FROM RecordType WHERE SobjectType = 'Account' AND DeveloperName = 'HH_Account' LIMIT 1].Id; - Account enriched = new Account( - Name = 'Optional net worth unavailable', RecordTypeId = householdType, - BillingStreet = '1 Test St', BillingCity = 'Powell River', BillingCountryCode = 'CA', - BillingStateCode = 'BC', BillingPostalCode = 'V0N 1T0', PRIZM_Score__c = 66, - Average_Income_in_Postal_Code__c = '$0 to $110,593', - Property_ownership_and_value__c = 'None or Unverifiable' - ); - insert enriched; - System.assertEquals(0, WCPrizmDailyEnrichment.selectPostalCodes().size(), - 'An assigned household must not be selected solely because optional net worth is unavailable.'); - Account pending = new Account( - Name = 'New postal code', RecordTypeId = householdType, - BillingStreet = '2 Test St', BillingCity = 'Powell River', BillingCountryCode = 'CA', - BillingStateCode = 'BC', BillingPostalCode = 'V8A 2P4' - ); - insert pending; - System.assertEquals(1, WCPrizmDailyEnrichment.selectPostalCodes().size()); + static void exhaustedCohortDoesNotFallBackToSalesforce() { + Test.setMock(HttpCalloutMock.class, new EmptyCohortMock()); + Test.startTest(); + WCPrizmDailyEnrichment.runDaily(); + System.assertEquals(1, Limits.getCallouts()); + System.assertEquals(0, Limits.getDmlStatements()); + Test.stopTest(); } private class MissingNetWorthMock implements HttpCalloutMock { diff --git a/test_cohort_queue.py b/test_cohort_queue.py new file mode 100644 index 0000000..8c671ad --- /dev/null +++ b/test_cohort_queue.py @@ -0,0 +1,98 @@ +import io +import os +import tempfile +import unittest +from concurrent.futures import ThreadPoolExecutor +from unittest.mock import patch + +from cache_manager_new import CacheManager +from cohort_queue import CohortQueue, postal_code +from app import app + + +class CohortTest(unittest.TestCase): + def setUp(self): + self.tmp = tempfile.TemporaryDirectory() + self.addCleanup(self.tmp.cleanup) + self.cache = CacheManager(os.path.join(self.tmp.name, 'test.db')) + self.queue = CohortQueue(self.cache.db_path) + self.codes = [f'V8A {n}A{i}' for n in range(3) for i in range(10)] + self.csv = 'Billing Country Code,Billing Zip/Postal Code\n' + ''.join('CA,' + p + '\n' for p in self.codes) + self.queue.import_accounts(io.StringIO(self.csv)) + + def test_normalization_import_and_exclusions(self): + self.assertIsNone(postal_code('D1D 1D1')) + self.queue.import_accounts(io.StringIO('Billing Country Code,Billing Zip/Postal Code\nCA,v8a0a0\nCA,V8A 0A0\nUS,90210\nCA,\n')) + summary = self.queue.summary() + self.assertEqual((summary['source_rows'], summary['excluded_rows'], summary['unique_codes'], summary['eligible_accounts']), (4, 2, 1, 2)) + with self.assertRaises(ValueError): + self.queue.import_accounts(io.StringIO('wrong\ncolumn\n')) + self.assertEqual(self.queue.summary()['unique_codes'], 1) + + def test_concurrency_limit_restart_next_day_and_no_repeat(self): + with ThreadPoolExecutor(max_workers=12) as pool: + claimed = list(pool.map(lambda _: self.queue.claim('2026-09-06'), range(20))) + codes = [p for p in claimed if p] + self.assertEqual(len(codes), 10) + self.assertEqual(len(set(codes)), 10) + for code in codes: + self.queue.finish(code, {'status': 'error'}) + queue = CohortQueue(self.cache.db_path) + self.assertIsNone(queue.claim('2026-09-06')) + self.assertNotIn(queue.claim('2026-09-07'), codes) + self.assertEqual(queue.summary()['failed'], 10) + + def test_captured_history_survives_expiry_deletion_and_reimport(self): + code = self.codes[0] + self.cache.cache_data(code, {'status': 'success', 'segment_number': '21'}, custom_duration_days=-1) + queue = CohortQueue(self.cache.db_path) + self.cache.delete_cached_data(code) + queue.import_accounts(io.StringIO(self.csv)) + self.assertEqual(queue.summary()['captured'], 1) + self.assertNotEqual(queue.claim('2026-09-06'), code) + + def test_api_auth_daily_limit_and_success_exclusion(self): + with patch('app.cache_manager', self.cache), patch.dict(os.environ, {'PRIZM_API_KEY': 'test-key', 'DASHBOARD_PASSWORD': 'dashboard'}), patch('app.get_prizm_code', side_effect=lambda code, **kw: {'postal_code': code, 'status': 'success'}) as lookup: + client = app.test_client() + self.assertEqual(client.post('/api/cohort/capture').status_code, 401) + first = client.post('/api/cohort/capture', headers={'X-API-Key': 'test-key'}) + self.assertEqual(first.status_code, 200) + self.assertEqual(len(first.json['results']), 1) + for _ in range(9): + client.post('/api/cohort/capture', headers={'X-API-Key': 'test-key'}) + second = client.post('/api/cohort/capture', headers={'X-API-Key': 'test-key'}) + self.assertEqual(second.json['results'], []) + self.assertEqual(lookup.call_count, 10) + self.assertEqual(self.queue.summary()['captured'], 10) + + def test_import_and_capture_require_api_key_even_with_dashboard_login(self): + import base64 + dashboard = 'Basic ' + base64.b64encode(b'dena:dashboard').decode() + with patch('app.cache_manager', self.cache), patch.dict(os.environ, {'PRIZM_API_KEY': 'key', 'DASHBOARD_USERNAME': 'dena', 'DASHBOARD_PASSWORD': 'dashboard'}): + client = app.test_client() + self.assertEqual(client.get('/api/cohort', headers={'Authorization': dashboard}).status_code, 200) + for path in ('/api/cohort/capture', '/api/cohort/import'): + self.assertEqual(client.post(path, headers={'Authorization': dashboard}).status_code, 401) + response = client.post('/api/cohort/import', data='wrong\ncolumn\n', headers={'X-API-Key': 'key'}) + self.assertEqual(response.status_code, 400) + self.assertEqual(self.queue.summary()['unique_codes'], 30) + + def test_capture_before_import_fails_closed(self): + with self.queue.connect() as db: + db.execute('DELETE FROM cohort_metadata') + with patch('app.cache_manager', self.cache), patch.dict(os.environ, {'PRIZM_API_KEY': 'key'}), patch('app.get_prizm_code') as lookup: + response = app.test_client().post('/api/cohort/capture', headers={'X-API-Key': 'key'}) + self.assertEqual(response.status_code, 503) + lookup.assert_not_called() + + def test_failed_and_interrupted_codes_not_repeated(self): + first = self.queue.claim('2026-09-06') + second = self.queue.claim('2026-09-06') + self.queue.finish(second, {'status': 'error', 'retryable': True}) + new = self.queue.claim('2026-09-07') + self.assertNotIn(new, [first, second]) + self.assertEqual(self.queue.summary()['incomplete'], 2) + + +if __name__ == '__main__': + unittest.main() From 78a2375d6f1f39f5f7aac4eb76b8ccbc1bd3e28b Mon Sep 17 00:00:00 2001 From: heybeaux Date: Sun, 6 Sep 2026 17:04:29 -0700 Subject: [PATCH 3/5] feat: merge historical PRIZM data while preserving current records --- .github/workflows/ci.yml | 4 +- app.py | 17 +++++++ cache_manager_new.py | 7 ++- docs/2026-09-06-upgrade.md | 42 +++++++++++++++++ historical_import.py | 94 ++++++++++++++++++++++++++++++++++++++ test_historical_import.py | 67 +++++++++++++++++++++++++++ 6 files changed, 228 insertions(+), 3 deletions(-) create mode 100644 historical_import.py create mode 100644 test_historical_import.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 04ab3b9..f07c73a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -30,10 +30,10 @@ jobs: pip install -r requirements.txt ruff - name: Lint - run: ruff check app.py prizm_client.py cache_manager_new.py cache_cli.py cohort_queue.py test_prizm_api.py test_cohort_queue.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 test_cohort_queue.py + run: python -m unittest test_prizm_api.py test_cohort_queue.py test_historical_import.py node: name: Node tests and build diff --git a/app.py b/app.py index 21e711c..0100a37 100644 --- a/app.py +++ b/app.py @@ -14,6 +14,7 @@ from cache_manager_new import cache_manager from cohort_queue import CohortQueue +from historical_import import merge_history from prizm_client import PrizmClient, PrizmLookupError, normalize_postal_code from segment_net_worth import average_household_net_worth, average_household_net_worth_amount @@ -438,6 +439,22 @@ def cohort_export(): return response +@app.route("/api/cache/import-history", methods=["POST"]) +def import_historical_cache(): + if not os.environ.get("PRIZM_API_KEY") or not has_valid_api_key(): + return jsonify({"error": "Historical import requires an API key"}), 401 + if request.content_length is None or request.content_length > 2_000_000: + return jsonify({"error": "A historical JSON payload under 2 MB is required"}), 400 + data = request.get_json(silent=True) + if not isinstance(data, dict): + return jsonify({"error": "Expected a JSON object"}), 400 + try: + result = merge_history(cache_manager.db_path, data.get("rows")) + except (ValueError, TypeError): + return jsonify({"error": "Invalid historical data; import rolled back"}), 400 + return jsonify(result) + + @app.route("/api/cohort/import", methods=["POST"]) def import_cohort(): if not os.environ.get("PRIZM_API_KEY") or not has_valid_api_key(): diff --git a/cache_manager_new.py b/cache_manager_new.py index 7b305e9..1b64abe 100644 --- a/cache_manager_new.py +++ b/cache_manager_new.py @@ -220,6 +220,7 @@ def _init_database(self): cursor, "postal_code_cache", { + "historical_source": "TEXT", "segment_number": "TEXT", "segment_name": "TEXT", "segment_description": "TEXT", @@ -709,8 +710,11 @@ def _dashboard_row(self, row: sqlite3.Row) -> Dict[str, Any]: 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")) - data["average_household_net_worth_amount"] = amount + 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"], @@ -852,6 +856,7 @@ def export_cache_csv(self, include_expired: bool = False) -> str: "average_household_net_worth", "average_household_net_worth_amount", "net_worth_source", + "data_source", "education", "urbanity", "occupation", diff --git a/docs/2026-09-06-upgrade.md b/docs/2026-09-06-upgrade.md index 833558a..640c2e1 100644 --- a/docs/2026-09-06-upgrade.md +++ b/docs/2026-09-06-upgrade.md @@ -76,3 +76,45 @@ The confirmed recent stall is 70 cache hits and zero upstream attempts August 31 September 6; latest cache timestamp is July 3. Full historical causality for every day since July 3 has not been independently reconstructed. Database SSH was unavailable with the configured local key; no new key was registered. + + +## Late-2025 historical dataset merge (prepared September 6) + +Beaux supplied `/Users/beauxwalton/prizm-api/prizm_cache_v2 copy.db` and requires +current Railway records to take priority over historical rows. Source database +passes SQLite quick_check; SHA-256 remains unchanged after analysis/import. + +- Source: 551 rows (517 success, 33 error, 1 invalid). Four rows have invalid + Canadian postal-code formats and are excluded with a separate audit CSV. +- 547 valid distinct codes: 159 overlap current data and remain entirely untouched. + 388 additional records (358 successes, 30 failures) are inserted. +- Preview has 551 total records: all 163 existing rows verified field-for-field + unchanged, plus 388 imports. The preview is built from an API-export fixture; + it is not a byte-for-byte production backup and must not replace the live DB. +- Cohort coverage increases from 106 to 362 of 484 codes; 122 remain uncaptured. + At ten successful new captures per day this is approximately 13 collection days. + +`historical_import.py` reads the source read-only and inserts only absent normalized +codes within a transaction. The authenticated `/api/cache/import-history` endpoint +applies the same rules to the private prepared payload during rollout. It checks +live overlaps at execution time, so intervening current captures also win. Reruns +are idempotent. Historical successes populate permanent capture history; failed +historical lookups do not count as captured. + +Preserve original expiry, confirmation and demographic values. The source has no +cached_at/scrape timestamp: imported cached_at stays NULL and unknown capture dates +stay blank. Do not infer scrape time from expiry or set it to import time. Historical +provenance is visible in details/CSV. All original expiry dates are already past; +imports remain available in full-cache views/exports but do not become fresh cache +hits. Net-worth amounts from the old scrape are retained as historical data rather +than overwritten by the segment fallback. + +Private artifacts: historical-import-payload.json, historical-merge-summary.json, +historical-rejected-postal-codes.csv and merged-history-preview.db under +`private/prizm-major-donors/` outside the repo. Do not commit the donor membership +list, database or import payload. + +Revised rollout: back up production; deploy updated Railway code; import historical +payload; verify current rows unchanged and historical counts; import cohort; verify +362 captured / 122 pending; then deploy already validated Salesforce classes. +No upstream scrape, Salesforce mutation or production database import was performed. diff --git a/historical_import.py b/historical_import.py new file mode 100644 index 0000000..8e356d1 --- /dev/null +++ b/historical_import.py @@ -0,0 +1,94 @@ +"""Import legacy cache rows without replacing current records or refreshing their age.""" +import argparse +import hashlib +import json +import sqlite3 +from pathlib import Path + +from cohort_queue import postal_code + +FIELDS = ( + 'postal_code', 'segment_number', 'segment_name', 'segment_description', 'who_they_are', + 'average_household_income', 'education', 'urbanity', 'average_household_net_worth', + 'occupation', 'diversity', 'family_life', 'tenure', 'home_type', 'status', 'confirmed', + 'expires_at', +) + + +def read_legacy(path): + """Open the source read-only; no source migration or modification.""" + with sqlite3.connect(Path(path).resolve().as_uri() + '?mode=ro', uri=True) as db: + db.row_factory = sqlite3.Row + if db.execute('PRAGMA quick_check').fetchone()[0] != 'ok': + raise ValueError('Source database integrity check failed') + rows = [dict(row) for row in db.execute('SELECT * FROM postal_code_cache ORDER BY postal_code')] + return rows + + +def merge_history(db_path, rows, source='late-2025'): + if not isinstance(rows, list) or len(rows) > 1000: + raise ValueError('Expected at most 1000 historical rows') + prepared, skipped, seen = [], 0, set() + for row in rows: + if not isinstance(row, dict): + raise ValueError('Historical rows must be objects') + code = postal_code(row.get('postal_code')) + if code is None: + skipped += 1 + continue + if code in seen: + raise ValueError('Duplicate normalized historical postal code') + seen.add(code) + if row.get('status') not in ('success', 'error', 'invalid') or not row.get('expires_at'): + raise ValueError('Historical status and original expiry are required') + values = {key: row.get(key) for key in FIELDS} + values['postal_code'] = code + prepared.append(values) + db = sqlite3.connect(db_path, timeout=30) + try: + with db: + db.execute('BEGIN IMMEDIATE') + db.execute('''CREATE TABLE IF NOT EXISTS historical_imports ( + postal_code TEXT PRIMARY KEY, source TEXT NOT NULL, + original_expires_at TEXT, imported_at TEXT DEFAULT CURRENT_TIMESTAMP + )''') + # Match codes across spacing/case variations, preserving every current field. + current = {postal_code(row[0]) for row in db.execute('SELECT postal_code FROM postal_code_cache')} + inserted = preserved = successes = 0 + for row in prepared: + code = row['postal_code'] + if row['status'] == 'success': + # Empty date explicitly means capture time unknown; never substitute import time. + db.execute("INSERT OR IGNORE INTO capture_history VALUES (?, '')", (code,)) + if code in current: + preserved += 1 + continue + columns = FIELDS + ('cached_at', 'average_household_net_worth_amount', 'historical_source') + amount = row.get('average_household_net_worth') + if not isinstance(amount, (int, float)): + amount = None + db.execute('INSERT INTO postal_code_cache (' + ','.join(columns) + ') VALUES (' + ','.join('?' for _ in columns) + ')', + [row[key] for key in FIELDS] + [None, amount, source]) + db.execute('INSERT OR IGNORE INTO historical_imports (postal_code,source,original_expires_at) VALUES (?,?,?)', + (code, source, row['expires_at'])) + current.add(code) + inserted += 1 + successes += row['status'] == 'success' + return {'source_rows': len(rows), 'inserted': inserted, 'inserted_successes': successes, + 'preserved_current': preserved, 'skipped_invalid_format': skipped} + finally: + db.close() + + +if __name__ == '__main__': + from cache_manager_new import CacheManager + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument('--source', required=True) + parser.add_argument('--target', required=True) + args = parser.parse_args() + if Path(args.source).resolve() == Path(args.target).resolve(): + parser.error('Source and target must differ') + rows = read_legacy(args.source) + CacheManager(args.target) + source_hash = hashlib.sha256(Path(args.source).read_bytes()).hexdigest() + print(json.dumps(merge_history(args.target, rows, 'late-2025 sha256:' + source_hash), indent=2)) diff --git a/test_historical_import.py b/test_historical_import.py new file mode 100644 index 0000000..f9a8900 --- /dev/null +++ b/test_historical_import.py @@ -0,0 +1,67 @@ +import os +import tempfile +import unittest +from unittest.mock import patch + +from app import app +from cache_manager_new import CacheManager +from cohort_queue import CohortQueue +from historical_import import merge_history + + +class HistoricalImportTest(unittest.TestCase): + def setUp(self): + self.tmp = tempfile.TemporaryDirectory() + self.addCleanup(self.tmp.cleanup) + self.cache = CacheManager(os.path.join(self.tmp.name, 'cache.db')) + self.old = {'postal_code':'v8a0a8', 'segment_number':'21', 'segment_name':'Historical', + 'status':'success', 'confirmed':0, 'expires_at':'2025-12-01 00:00:00', + 'average_household_net_worth':123456, 'home_type':'Detached'} + + def test_current_row_preserved_entirely_and_normalized_overlap(self): + self.cache.cache_data('V8A 0A8', {'status':'error', 'message':'Current result'}) + with self.cache._connect() as db: + before = tuple(db.execute('SELECT * FROM postal_code_cache').fetchone()) + result = merge_history(self.cache.db_path, [self.old]) + with self.cache._connect() as db: + self.assertEqual(before, tuple(db.execute('SELECT * FROM postal_code_cache').fetchone())) + self.assertEqual(result['preserved_current'], 1) + self.assertEqual(result['inserted'], 0) + + def test_import_preserves_age_values_and_history_after_cache_deletion(self): + result = merge_history(self.cache.db_path, [self.old]) + self.assertEqual(result['inserted'], 1) + self.assertIsNone(self.cache.get_cached_data('V8A 0A8')) + row = self.cache.list_cache_entries(include_expired=True)[0] + self.assertIsNone(row['cached_at']) + self.assertEqual(row['expires_at'], self.old['expires_at']) + self.assertEqual(row['average_household_net_worth_amount'], 123456) + self.assertIn('Historical import', row['data_source']) + self.assertEqual(merge_history(self.cache.db_path, [self.old])['inserted'], 0) + self.cache.delete_cached_data('V8A 0A8') + queue = CohortQueue(self.cache.db_path) + with queue.connect() as db: + db.execute("INSERT INTO donor_cohort VALUES ('V8A 0A8',1)") + self.assertIsNone(queue.claim('2026-09-07')) + self.assertEqual(queue.entries()[0]['captured_at'], '') + self.assertEqual(merge_history(self.cache.db_path, [self.old])['inserted'], 1) + self.assertIsNone(queue.claim('2026-09-08')) + + def test_bad_row_aborts_whole_import_and_bad_formats_are_reported(self): + with self.assertRaises(ValueError): + merge_history(self.cache.db_path, [self.old, {'postal_code':'V8A 1A1','status':'oops'}]) + self.assertEqual(self.cache.list_cache_entries(include_expired=True), []) + result = merge_history(self.cache.db_path, [{'postal_code':'90210','status':'invalid'}]) + self.assertEqual(result['skipped_invalid_format'], 1) + + def test_endpoint_auth_and_input_validation(self): + with patch('app.cache_manager', self.cache), patch.dict(os.environ, {'PRIZM_API_KEY':'key'}): + client = app.test_client() + self.assertEqual(client.post('/api/cache/import-history', json={'rows':[self.old]}).status_code, 401) + response = client.post('/api/cache/import-history', headers={'X-API-Key':'key'}, json={'rows':[self.old]}) + self.assertEqual(response.status_code, 200) + self.assertEqual(response.json['inserted'], 1) + + +if __name__ == '__main__': + unittest.main() From 27d6c26569b27a59cc43a6dfad81bc562bf40175 Mon Sep 17 00:00:00 2001 From: heybeaux Date: Sun, 6 Sep 2026 17:16:34 -0700 Subject: [PATCH 4/5] fix: snapshot production cache and serialize worker migrations --- .gitignore | 2 ++ cache_manager_new.py | 19 ++++++++++++++++++- test_historical_import.py | 10 ++++++++++ 3 files changed, 30 insertions(+), 1 deletion(-) diff --git a/.gitignore b/.gitignore index 102e7b1..b6d5a16 100644 --- a/.gitignore +++ b/.gitignore @@ -56,3 +56,5 @@ debug_screenshots/ missing_postal_codes.txt postal_code_export*.csv postal_code_v2_export*.csv + +*.migration.lock diff --git a/cache_manager_new.py b/cache_manager_new.py index 1b64abe..bb9d6f4 100644 --- a/cache_manager_new.py +++ b/cache_manager_new.py @@ -6,6 +6,7 @@ """ import csv +import fcntl import io import json import logging @@ -27,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): diff --git a/test_historical_import.py b/test_historical_import.py index f9a8900..60c4622 100644 --- a/test_historical_import.py +++ b/test_historical_import.py @@ -54,6 +54,16 @@ def test_bad_row_aborts_whole_import_and_bad_formats_are_reported(self): result = merge_history(self.cache.db_path, [{'postal_code':'90210','status':'invalid'}]) self.assertEqual(result['skipped_invalid_format'], 1) + def test_production_startup_snapshot_is_preserved(self): + import sqlite3 + self.cache.cache_data('V8A 0A8', {'status':'success', 'segment_number':'21'}) + with patch.dict(os.environ, {'RAILWAY_ENVIRONMENT_ID':'test-production'}): + CacheManager(self.cache.db_path) + self.cache.delete_cached_data('V8A 0A8') + CacheManager(self.cache.db_path) + with sqlite3.connect(self.cache.db_path + '.before-cohort-upgrade.db') as db: + self.assertEqual(db.execute('SELECT COUNT(*) FROM postal_code_cache').fetchone()[0], 1) + def test_endpoint_auth_and_input_validation(self): with patch('app.cache_manager', self.cache), patch.dict(os.environ, {'PRIZM_API_KEY':'key'}): client = app.test_client() From 71c6354c0c40edc6d90bc43c057b6c92c7686214 Mon Sep 17 00:00:00 2001 From: heybeaux Date: Sun, 6 Sep 2026 17:18:52 -0700 Subject: [PATCH 5/5] ci: pin Ruff to the validated lint version --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f07c73a..936f4a4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -27,7 +27,7 @@ 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 cohort_queue.py historical_import.py test_prizm_api.py test_cohort_queue.py test_historical_import.py