#!/usr/bin/env python3 """Intake watcher: FormSubmit quote-request email -> parse -> price -> proposal. Run by cron every 15 min. State in state.json (processed message ids). For each new "Solutions review request" email in sean@seanmattsson.com: 1. parse fields from the request_details block 2. save intake/intake_.json 3. price via price_proposal.price_intake 4. if ready: build PDF proposal -> proposals/ 5. update queue.json for the dashboard Prints a JSON summary for the cron report. """ import os, json, re, subprocess, sys, datetime BASE = os.path.dirname(os.path.abspath(__file__)) STATE = os.path.join(BASE, "state.json") QUEUE = os.path.join(BASE, "dashboard", "queue.json") OPPS = "/home/hatch/workspace/goals/sales-weapon-ae-copilot/hidden_files/prospect-queues/opportunities.json" DEPLOY_CONFIG = os.path.join(BASE, "deploy_config.json") ACCOUNT = "d48c9e2643f643aeb3c4dacf347245b8" # sean@seanmattsson.com def cli(*args): r = subprocess.run(["hatch_gws_cli", "gmail", "--account", ACCOUNT] + list(args), capture_output=True, text=True, timeout=120) return r.stdout def load_state(): if os.path.exists(STATE): return json.load(open(STATE)) return {"processed": []} SECTION_HEADERS = ["STRATIX SYSTEMS \u2014 SOLUTIONS REVIEW REQUEST", "CONTACT", "PRIORITIES", "PRINT & SCAN", "LEASE & IT", "CONFIGURATION"] DISCLAIMER = "This is a request for review and quotation" TEST_PAT = re.compile(r"^(test|ddd|asdf|example)\b|\bplease ignore\b|example\.com$", re.I) def parse_details_blob(blob, label_to_key): """Split the request_details text on known 'Label:' boundaries.""" labels = sorted(label_to_key.keys(), key=len, reverse=True) pat = re.compile("(" + "|".join(re.escape(l) for l in labels) + r")\s*:") matches = list(pat.finditer(blob)) fields = {} for i, m in enumerate(matches): key = label_to_key.get(m.group(1)) if not key: continue start = m.end() end = matches[i + 1].start() if i + 1 < len(matches) else len(blob) val = blob[start:end] for htxt in SECTION_HEADERS: val = val.replace(htxt, " ") if DISCLAIMER in val: val = val.split(DISCLAIMER)[0] val = re.sub(r"\s+", " ", val).strip(" -|") if val: fields[key] = val return fields def parse_email_body(text, label_to_key): """Extract fields from a FormSubmit table email.""" fields = {} for row_key, fkey in [("name", "contactName"), ("email", "email"), ("phone", "phone"), ("company", "company")]: m = re.search(r"\|\s*%s\s*\|\s*(.*?)\s*\|" % re.escape(row_key), text) if m and m.group(1).strip(): fields[fkey] = m.group(1).strip() m = re.search(r"request_details\s*\|\s*(.*?)\n\nSubmitted at", text, re.S) blob = m.group(1) if m else text fields.update(parse_details_blob(blob, label_to_key)) return fields def is_test_submission(fields): for k in ("company", "contactName", "email"): if fields.get(k) and TEST_PAT.search(fields[k]): return True return False def company_slug(company): s = re.sub(r"[^a-z0-9]+", "-", (company or "customer").lower()).strip("-") return (s or "customer")[:40] def lead_score(fields, priced): """HOT/WARM/COOL scoring so Sean's alerts lead with the money (council 2026-10-06).""" score, reasons = 0, [] timing = (fields.get("timing") or "").lower() if "as soon as possible" in timing: score += 2; reasons.append("ASAP timing") elif "within 30 days" in timing or "3 months" in timing: score += 1; reasons.append("near-term timing") try: ml = float(str(fields.get("monthsLeft") or "nan").replace(",", "")) if ml <= 6: score += 2; reasons.append("lease expires in %g mo" % ml) except Exception: pass if (fields.get("interests") or "").strip(): score += 1; reasons.append("add-on interests flagged") if (priced.get("total_monthly") or 0) >= 500: score += 1; reasons.append("large deal size") tier = "HOT" if score >= 3 else "WARM" if score >= 1 else "COOL" return {"tier": tier, "score": score, "reasons": reasons} def publish_interactive(slug, html_path): """Deploy the interactive proposal to its per-customer URL. Uses deploy_proposal.publish_interactive (Cloudflare Pages). Returns the live URL, or None on any failure (then the HTML file rides along as a draft attachment instead). Controlled by deploy_config.json. """ try: cfg = json.load(open(DEPLOY_CONFIG)) except Exception: cfg = {} if not cfg.get("enabled"): return None try: from deploy_proposal import publish_interactive as _publish return _publish(slug, html_path) except Exception: return None def upsert_opportunity(entry, priced, draft_id, live_url): """Write/update the PROPOSAL-stage opportunity so the follow-up engine picks it up.""" try: opps = json.load(open(OPPS)) except Exception: opps = [] cust = priced.get("customer", {}) account = cust.get("company") or "Unknown" email = (cust.get("email") or "").lower() today = datetime.date.today().isoformat() opp = None for o in opps: if (email and (o.get("email") or "").lower() == email) or \ (o.get("account") or "").lower() == account.lower(): opp = o break event = ("Quote request via site; proposal auto-built (%s/mo, %d-mo %s); " "customer draft %s for Sean's review%s." % ( "$%s" % format(priced.get("total_monthly", 0), ",.2f"), priced.get("term_months", 63), priced.get("lease_type", "FMV"), "created" if draft_id else "FAILED to create", "; live link %s" % live_url if live_url else "")) if opp is None: opp = { "id": "opp-%s-%s" % (company_slug(account), today.replace("-", "")), "account": account, "contact": cust.get("contact_name", ""), "email": cust.get("email", ""), "stage": "PROPOSAL", "value": priced.get("total_monthly"), "lane": "net_new", "source": "quote_site", "proposal_date": today, "proposal_url": live_url, "valid_until": priced.get("valid_until"), "next_action": "Sean reviews draft and sends proposal", "next_action_date": today, "verified": True, "history": [{"date": today, "event": event}], } opps.append(opp) else: opp["stage"] = "PROPOSAL" opp["proposal_date"] = today opp["proposal_url"] = live_url or opp.get("proposal_url") opp["valid_until"] = priced.get("valid_until") opp["value"] = priced.get("total_monthly") opp["next_action"] = "Sean reviews draft and sends proposal" opp["next_action_date"] = today opp.setdefault("history", []).append({"date": today, "event": event}) json.dump(opps, open(OPPS, "w"), indent=1) return opp["id"] def _fingerprint(fields): day = datetime.date.today().isoformat() return "%s|%s|%s" % ((fields.get("email") or "").lower().strip(), (fields.get("company") or "").lower().strip(), day) def _cf_api(token, method, path, data=None): import urllib.request as _url body = json.dumps(data).encode() if data is not None else None req = _url.Request("https://api.cloudflare.com/client/v4" + path, data=body, method=method, headers={"Authorization": "Bearer " + token, "Content-Type": "application/json"}) return json.load(_url.urlopen(req, timeout=30)) KV_NAMESPACE = "67116ce521ef4228ae4827ae5acd8885" KV_ACCOUNT = "76ef691dfd113325265b6a77d27ce435" KV_SAFETY_DELAY_MIN = 30 # give the FormSubmit email a head start def kv_safety_net(state, fingerprints, process_fields): """Backup intake: submissions the form logged to KV but whose FormSubmit email never arrived (delayed/dropped). Only entries older than KV_SAFETY_DELAY_MIN minutes with no matching processed fingerprint.""" new_entries = [] try: from deploy_proposal import find_deploy_token token = find_deploy_token() if not token: return new_entries now_ms = int(datetime.datetime.now().timestamp() * 1000) cutoff = now_ms - KV_SAFETY_DELAY_MIN * 60000 watermark = state.get("kv_watermark", 0) keys = _cf_api(token, "GET", "/accounts/%s/storage/kv/namespaces/%s/keys" % (KV_ACCOUNT, KV_NAMESPACE)) max_ts = watermark for k in keys.get("result", []): name = k.get("name", "") m = re.match(r"intake_(\d+)_[a-z0-9]+$", name) if not m: continue ts = int(m.group(1)) max_ts = max(max_ts, ts) if ts <= watermark or ts > cutoff: continue try: val = _cf_api(token, "GET", "/accounts/%s/storage/kv/namespaces/%s/values/%s" % (KV_ACCOUNT, KV_NAMESPACE, name)) except Exception: continue fields = val.get("fields", {}) if isinstance(val, dict) else {} # normalize checkbox arrays fields = {kk: (", ".join(vv) if isinstance(vv, list) else vv) for kk, vv in fields.items()} if is_test_submission(fields): continue fp = _fingerprint(fields) if fp in fingerprints: continue # the email already got processed entry = process_fields(fields, "kv_%s" % name.replace("intake_", ""), "kv_backup") entry["source"] = "kv_backup" new_entries.append(entry) fingerprints.add(fp) state["kv_watermark"] = max_ts except Exception as e: print("kv safety net error: %s" % str(e)[:200]) return new_entries def main(): label_to_key = json.load(open(os.path.join(BASE, "data", "label_to_key.json"))) state = load_state() processed = set(state["processed"]) fingerprints = set(state.get("fingerprints", [])) out = cli("+triage", "--query", 'subject:"Solutions review request"', "--max", "20", "--format", "json") try: msgs = json.loads(out) items = msgs.get("messages", msgs) if isinstance(msgs, dict) else msgs except Exception: items = [] if not isinstance(items, list): items = [] sys.path.insert(0, BASE) from price_proposal import price_intake from build_proposal import save_proposal_pdf from build_interactive import save_interactive from draft_builder import build_customer_draft, build_gap_draft, create_draft run_at = datetime.datetime.now().isoformat(timespec="seconds") def process_fields(fields, iid, source): """Full pipeline for one parsed intake: price -> build -> draft -> opp.""" customer = { "company": fields.get("company", ""), "contact_name": fields.get("contactName", fields.get("name", "")), "phone": fields.get("phone", ""), "email": fields.get("email", ""), "address": fields.get("businessLocation", ""), } intake = {"customer": customer, "model": fields.get("model", ""), "options": fields.get("options", ""), "bwVolume": fields.get("bwVolume", ""), "colorVolume": fields.get("colorVolume", ""), "monthlyPayment": fields.get("monthlyPayment", ""), "currentModel": fields.get("currentModel", ""), "raw": fields} slug = "%s-%s" % (company_slug(customer["company"]), datetime.date.today().strftime("%Y%m%d")) json.dump(intake, open(os.path.join(BASE, "intake", iid + ".json"), "w"), indent=1) priced = price_intake(intake) entry = {"id": iid, "slug": slug, "source": source, "customer": customer, "model": priced.get("model"), "status": priced["status"], "model_recommended": priced.get("model_recommended", False), "model_note": priced.get("model_note", ""), "gaps": priced.get("gaps", []), "review_flags": priced.get("review_flags", []), "lead": lead_score(fields, priced), "total_monthly": priced.get("total_monthly"), "lease_monthly": priced.get("lease_monthly"), "valid_until": priced.get("valid_until"), "built_at": run_at} if priced["status"] == "ready": pdf = save_proposal_pdf(priced) entry["proposal_pdf"] = os.path.relpath(pdf, BASE) ipath = save_interactive(priced) entry["interactive_path"] = os.path.relpath(ipath, BASE) live_url = publish_interactive(slug, ipath) entry["live_url"] = live_url subject, html_body = build_customer_draft(priced, live_url) attachments = [pdf] + ([] if live_url else [ipath]) draft_id = create_draft(customer["email"], subject, html_body, attachments) entry["draft_id"] = draft_id entry["draft_status"] = "ready_for_review" if draft_id else "draft_failed" entry["opp_id"] = upsert_opportunity(entry, priced, draft_id, live_url) else: # Gaps: no proposal. Draft a "quick question" email instead of stalling. subject, html_body = build_gap_draft(priced) draft_id = create_draft(customer["email"], subject, html_body, []) entry["draft_id"] = draft_id entry["draft_status"] = "gap_questions_ready" if draft_id else "draft_failed" return entry results = [] for m in items: mid = m.get("id") if not mid or mid in processed: continue body = cli("+read", "--id", mid, "--format", "json") try: msg = json.loads(body) text = msg.get("body_text") or msg.get("body") or body except Exception: text = body fields = parse_email_body(text if isinstance(text, str) else str(text), label_to_key) if is_test_submission(fields): processed.add(mid) results.append({"id": "skipped", "message_id": mid, "status": "skipped_test", "company": fields.get("company")}) continue entry = process_fields(fields, "qr_%s" % mid[-8:], "email") entry["message_id"] = mid fingerprints.add(_fingerprint(fields)) results.append(entry) processed.add(mid) # Safety net: KV-logged submissions whose email never arrived. results.extend(kv_safety_net(state, fingerprints, process_fields)) state["processed"] = sorted(processed) state["fingerprints"] = sorted(fingerprints) json.dump(state, open(STATE, "w"), indent=1) # merge into dashboard queue queue = json.load(open(QUEUE)) if os.path.exists(QUEUE) else [] have = {q["id"] for q in queue} for r in results: if r["id"] not in have and not r["id"].startswith("skipped"): queue.insert(0, r) json.dump(queue, open(QUEUE, "w"), indent=1) print(json.dumps({"new": len(results), "results": results}, indent=1)) if __name__ == "__main__": main()