"""Outreach — the internal campaign engine (the HubSpot-parity build, crm-outreach brief v1). Select a segment of customers → render a personalized template per recipient → queue one mail.mail each through the ONE sanctioned write path (modules/collections_send.py client: mail.mail create only, SAFE_MODE allow-list enforced in the data layer) → read outcomes back from mail.mail state → attribute orders placed within the window to the campaign (the thing HubSpot can't do: revenue attribution straight off sale.order). Store keys: suppression {email_lower: {reason: unsub|bounce|manual, source, at}} outreach_templates {id: {name, class: marketing|transactional, subject, body_html, ...}} campaigns {seq, items: {id: {name, template_id, segment, audience[], excluded[], status draft|queued, scheduled_base, send_log[], sent_at, created_by}}} Compliance (CAN-SPAM — applies fully to B2B): - marketing-class sends get the physical-address + unsubscribe footer appended ALWAYS; - the suppression list is enforced HERE (same choke point as SAFE_MODE), not in the UI; - unsubscribe v1 = reply-to-unsubscribe (explicitly valid) + List-Unsubscribe mailto header; - frequency governor: max 1 marketing email per dealer per 7 days across ALL campaigns. O365 limits: 30 msgs/min → sends are staggered via mail.mail.scheduled_date in 1-min buckets of 25; 10k recipients/day is far above the dealer book. """ import os import html as _html import datetime as dt from jinja2 import Environment, StrictUndefined import core.store as store import core.odoo as O import core.periods as P import modules.customers as cust import modules.collections_send as cs K_SUPP = 'suppression' K_TPL = 'outreach_templates' K_CAMP = 'campaigns' BATCH_PER_MINUTE = 25 # O365 message rate limit is 30/min — stay under it GOVERNOR_DAYS = 7 # max 1 marketing email per dealer per 7 days ATTRIB_DAYS = 14 # orders within N days of send count toward the campaign _env = Environment(undefined=StrictUndefined, autoescape=False) def _now(): return dt.datetime.now().strftime('%Y-%m-%d %H:%M') # ------------------------------------------------------------------ suppression def suppression(): return store.get(K_SUPP) or {} def is_suppressed(email): return str(email or '').strip().lower() in suppression() def suppress(email, reason, source, by='?'): e = str(email or '').strip().lower() if not e: return def _add(d): d[e] = {'reason': reason, 'source': source, 'by': by, 'at': _now()} return d store.update(K_SUPP, _add) def unsuppress(email): e = str(email or '').strip().lower() def _rm(d): d.pop(e, None) return d store.update(K_SUPP, _rm) # ------------------------------------------------------------------ templates DEFAULT_TEMPLATES = { 'reorder-v1': { 'name': 'Reorder reminder', 'class': 'marketing', 'subject': 'Time to restock, {{ customer }}?', 'body_html': ( '

Hi {{ customer }},

' '

It has been {{ days_since }} days since your last order with us' '{% if last_order %} (placed {{ last_order }}){% endif %} — based on your usual ' 'rhythm, you may be running low.

' '{% if top_skus %}

Your recent favorites:

{% endif %}' '

Reply to this email or reach out to {{ agent }} and we will get your next ' 'order moving.

' '

Thank you for your business,
{{ company }}

'), }, 'winback-v1': { 'name': 'Win-back', 'class': 'marketing', 'subject': 'We miss you at {{ company }}, {{ customer }}', 'body_html': ( '

Hi {{ customer }},

' '

We noticed your orders have slowed this year and wanted to check in — our new ' 'catalog has landed and we would love to get you back to the table.

' '

Reply to this email or contact {{ agent }} to hear what is new.

' '

Thank you,
{{ company }}

'), }, } def templates(): t = store.get(K_TPL) or {} if not t: stamped = {k: {**v, 'created_by': 'system', 'at': _now()} for k, v in DEFAULT_TEMPLATES.items()} try: store.update(K_TPL, lambda d: (d.update(stamped), d)[1] if not d else d) t = store.get(K_TPL) or stamped except Exception: t = stamped return t def save_template(tid, name, tclass, subject, body_html, by='?'): tid = (tid or name).strip().lower().replace(' ', '-') def _set(d): d[tid] = {'name': name, 'class': tclass, 'subject': subject, 'body_html': body_html, 'created_by': by, 'at': _now()} return d store.update(K_TPL, _set) return tid def render(subject_tmpl, body_tmpl, ctx): """StrictUndefined — a missing token fails the PREVIEW, never a live send.""" return (_env.from_string(subject_tmpl).render(**ctx), _env.from_string(body_tmpl).render(**ctx)) # ------------------------------------------------------------------ company footer (CAN-SPAM) _company_cache = {} def company_address(): """Physical postal address for the compliance footer. Env override, else res.company.""" if 'addr' in _company_cache: return _company_cache['addr'] addr = os.environ.get('COMPANY_ADDRESS', '').strip() if not addr: try: rows = O.search_read('res.company', [], ['name', 'street', 'street2', 'city', 'state_id', 'zip']) # multi-company instance: take the first company with a real street address c = next((r for r in rows if r.get('street')), None) if c: bits = [c.get('street') or '', c.get('street2') or '', c.get('city') or '', O.m2o_name(c.get('state_id')), c.get('zip') or ''] addr = ', '.join(b for b in bits if b) except Exception: addr = '' _company_cache['addr'] = addr return addr def marketing_footer(): addr = company_address() return (f"
{_html.escape(cs.COMPANY)}" + (f" · {_html.escape(addr)}" if addr else '') + "
You are receiving this because you are a wholesale customer of " f"{_html.escape(cs.COMPANY)}. To stop receiving these emails, reply with " "the word unsubscribe and we will remove you within 10 business days." "
") def wrap_body(inner_html, tclass): shell = (f"
{inner_html}" + (marketing_footer() if tclass == 'marketing' else '') + '
') return shell # ------------------------------------------------------------------ audiences SEGMENT_SOURCES = { 'reorder': 'Reorder due — overdue vs their own cadence', 'winback': 'Win-back — bought LY, down or gone this year', 'new': 'New / reactivated this year', } def build_audience(source, team_id=None, agent=None, limit=500): """Segment spec → candidate rows [{pid, customer, metric, agent?}] (active-list semantics: evaluated NOW; the campaign freezes a snapshot).""" pids_scope = cust.agent_partner_ids(agent) if agent else None if source == 'reorder': rows = cust.contact_recommendations(team_id=team_id, limit=limit, agent_pids=pids_scope) return [{'pid': r['pid'], 'customer': r['customer'], 'agent': r.get('agent', ''), 'metric': r.get('est_missed', 0.0), 'metric_label': 'est. missed $'} for r in rows] if source == 'winback': rows = cust.at_risk(team_id=team_id, limit=limit, agent_pids=pids_scope) return [{'pid': r['pid'], 'customer': r['customer'], 'agent': '', 'metric': r.get('at_risk', 0.0), 'metric_label': 'at-risk $'} for r in rows] if source == 'new': rows = cust.new_customers(team_id=team_id, limit=limit, agent_pids=pids_scope) return [{'pid': r['pid'], 'customer': r['customer'], 'agent': '', 'metric': r.get('rev_ytd', 0.0), 'metric_label': 'YTD $'} for r in rows] raise ValueError(f'unknown segment source {source}') def _partner_info(pids): """{pid: {email, name, credit}} in one read.""" if not pids: return {} rows = O.search_read('res.partner', [('id', 'in', list(pids))], ['name', 'email', 'credit']) return {r['id']: r for r in rows} def _tokens_bulk(pids, team_id=None): """Per-recipient template context, pulled in TWO read_groups (never per-recipient calls): cadence (last order / days since) + top-3 SKUs by LTM revenue.""" cad = cust._cadence_bulk(team_id=team_id) lf, lt = P.ltm(P.today()) # ISO strings already (periods convention) g = O.read_group('sale.order.line', O.sale_line_domain(lf, lt, team_id, partner_ids=list(pids)), ['price_subtotal:sum'], ['order_partner_id', 'product_id'], lazy=False) by_pid = {} for r in g: pid = O.m2o_id(r.get('order_partner_id')) name = O.m2o_name(r.get('product_id')) if pid and name: by_pid.setdefault(pid, []).append((r.get('price_subtotal') or 0.0, name)) top = {pid: [n for _, n in sorted(v, reverse=True)[:3]] for pid, v in by_pid.items()} out = {} for pid in pids: c = cad.get(pid, {}) out[pid] = {'last_order': c.get('last_order') or '', 'days_since': int(c.get('days_since') or 0), 'top_skus': top.get(pid, [])} return out def _agent_disp(a): """Agent token for templates — the '(none)' placeholder must never reach a customer.""" a = str(a or '').strip() return a if a and a.lower() != '(none)' else 'your account manager' def _recent_marketing_sends(): """{email_lower: last marketing send date iso} across all campaigns — the frequency governor.""" out = {} camps = (store.get(K_CAMP) or {}).get('items', {}) for c in camps.values(): if c.get('template_class') != 'marketing': continue day = (c.get('sent_at') or '')[:10] if not day: continue for s in c.get('send_log', []): e = str(s.get('email', '')).lower() if e and (e not in out or out[e] < day): out[e] = day return out # ------------------------------------------------------------------ campaigns def campaigns(): d = store.get(K_CAMP) or {} return d.get('items', {}) def create_campaign(name, template_id, segment_spec, audience_rows, by='?'): """Freeze the audience snapshot with exclusions applied UP FRONT (each with its reason — the excluded rows are shown, never silently dropped): no email / suppressed / marketing governor (emailed < GOVERNOR_DAYS ago) / duplicate email within the audience.""" tpl = templates().get(template_id) or {} tclass = tpl.get('class', 'marketing') info = _partner_info([r['pid'] for r in audience_rows]) recent = _recent_marketing_sends() if tclass == 'marketing' else {} supp = suppression() cutoff = (dt.date.today() - dt.timedelta(days=GOVERNOR_DAYS)).isoformat() seen = set() audience, excluded = [], [] for r in audience_rows: p = info.get(r['pid']) or {} email = str(p.get('email') or '').strip() el = email.lower() if not email: excluded.append({**r, 'reason': 'no email address'}) elif el in supp: excluded.append({**r, 'reason': f"suppressed ({supp[el].get('reason', '?')})"}) elif tclass == 'marketing' and recent.get(el, '') >= cutoff: excluded.append({**r, 'reason': f'emailed within {GOVERNOR_DAYS}d (governor)'}) elif el in seen: excluded.append({**r, 'reason': 'duplicate email in audience'}) else: seen.add(el) audience.append({**r, 'email': email}) created = {} def _add(d): d.setdefault('seq', 0) d.setdefault('items', {}) d['seq'] += 1 cid = f"c{d['seq']}" d['items'][cid] = { 'id': cid, 'name': name, 'template_id': template_id, 'template_class': tclass, 'segment': segment_spec, 'audience': audience, 'excluded': excluded, 'status': 'draft', 'send_log': [], 'sent_at': None, 'created_by': by, 'created_at': _now(), } created['id'] = cid return d store.update(K_CAMP, _add) return created['id'], len(audience), len(excluded) def preview_samples(campaign, n=3, team_id=None): """Render the template for the first n recipients — StrictUndefined surfaces bad tokens here, at preview time.""" tpl = templates().get(campaign['template_id']) or {} aud = campaign['audience'][:n] toks = _tokens_bulk([r['pid'] for r in aud], team_id) out = [] for r in aud: ctx = {'customer': r['customer'], 'agent': _agent_disp(r.get('agent')), 'company': cs.COMPANY, 'open_balance': 0.0, **toks.get(r['pid'], {})} subj, body = render(tpl.get('subject', ''), tpl.get('body_html', ''), ctx) out.append({'to': r['email'], 'subject': subj, 'body': wrap_body(body, tpl.get('class', 'marketing'))}) return out def send_campaign(cid, base_dt_utc=None, by='?', team_id=None): """Queue every audience mail (personalized, one mail.mail per recipient, staggered scheduled_date ≤ BATCH_PER_MINUTE/min). SAFE_MODE is checked per recipient HERE: while the guardrail is ON, non-allow-listed recipients are SKIPPED (logged), so a test send to yourself works without touching real customers. Returns (queued, skipped).""" camps = campaigns() c = camps.get(cid) if not c or c.get('status') != 'draft': raise ValueError('campaign missing or already sent') tpl = templates().get(c['template_id']) or {} tclass = tpl.get('class', 'marketing') base = base_dt_utc or dt.datetime.utcnow() toks = _tokens_bulk([r['pid'] for r in c['audience']], team_id) info = _partner_info([r['pid'] for r in c['audience']]) odoo = cs.Odoo() log = [] queued = skipped = 0 for i, r in enumerate(c['audience']): email = r['email'] if is_suppressed(email): # re-check at send time log.append({**r, 'result': 'suppressed'}) skipped += 1 continue if not cs.safe_recipient_ok(email): log.append({**r, 'result': 'safe-mode-skip'}) skipped += 1 continue ctx = {'customer': r['customer'], 'agent': _agent_disp(r.get('agent')), 'company': cs.COMPANY, 'open_balance': float((info.get(r['pid']) or {}).get('credit') or 0.0), **toks.get(r['pid'], {})} try: subj, body = render(tpl.get('subject', ''), tpl.get('body_html', ''), ctx) when = base + dt.timedelta(minutes=queued // BATCH_PER_MINUTE) payload = { 'subject': subj, 'body_html': wrap_body(body, tclass), 'email_to': email, 'email_from': cs.SENDER_DISPLAY, 'reply_to': cs.REPLY_TO, 'mail_server_id': cs.ROYAL_MAIL_SERVER_ID, 'author_id': cs.ROYAL_AUTHOR_ID, 'state': 'outgoing', 'auto_delete': False, 'scheduled_date': when.strftime('%Y-%m-%d %H:%M:%S'), 'model': 'res.partner', 'res_id': r['pid'], # audit trail in chatter 'headers': repr({'List-Unsubscribe': f''}), } mid = odoo.queue_mail(payload) log.append({**r, 'result': 'queued', 'mail_id': mid, 'scheduled': when.strftime('%H:%M')}) queued += 1 except Exception as e: log.append({**r, 'result': f'error: {str(e)[:90]}'}) skipped += 1 def _fin(d): item = d.get('items', {}).get(cid) if item: item['status'] = 'queued' item['send_log'] = log item['sent_at'] = _now() item['sent_by'] = by return d store.update(K_CAMP, _fin) return queued, skipped def refresh_outcomes(cid): """Read mail.mail state back for a campaign's queued mails (sent / exception + reason).""" c = campaigns().get(cid) or {} ids = [s['mail_id'] for s in c.get('send_log', []) if s.get('mail_id')] if not ids: return {} rows = O.search_read('mail.mail', [('id', 'in', ids)], ['state', 'failure_reason']) by_id = {r['id']: r for r in rows} def _upd(d): item = d.get('items', {}).get(cid) if item: for s in item.get('send_log', []): r = by_id.get(s.get('mail_id')) if r: s['state'] = r.get('state') fr = r.get('failure_reason') if fr: s['failure'] = str(fr)[:120] return d store.update(K_CAMP, _upd) return by_id def attribution(cid, days=ATTRIB_DAYS): """Orders placed by audience customers within `days` of send → per-customer rows + totals. Correlation window, honestly labeled — not causal ML.""" c = campaigns().get(cid) or {} sent = (c.get('sent_at') or '')[:10] if not sent: return {'rows': [], 'orders': 0, 'revenue': 0.0, 'window': days, 'from': None, 'to': None} pids = [r['pid'] for r in c.get('audience', [])] if not pids: return {'rows': [], 'orders': 0, 'revenue': 0.0, 'window': days, 'from': sent, 'to': None} end = (dt.date.fromisoformat(sent) + dt.timedelta(days=days)).isoformat() g = O.read_group('sale.order', [('partner_id', 'in', pids), ('state', 'in', ['sale', 'done']), ('date_order', '>=', f'{sent} 00:00:00'), ('date_order', '<=', f'{end} 23:59:59')], ['amount_untaxed:sum'], ['partner_id'], lazy=False) names = {r['pid']: r['customer'] for r in c.get('audience', [])} rows = [] for r in g: pid = O.m2o_id(r.get('partner_id')) if not pid: continue rows.append({'pid': pid, 'customer': names.get(pid) or O.m2o_name(r.get('partner_id')), 'orders': r.get('__count') or 0, 'revenue': r.get('amount_untaxed') or 0.0}) rows.sort(key=lambda x: -x['revenue']) return {'rows': rows, 'orders': sum(x['orders'] for x in rows), 'revenue': sum(x['revenue'] for x in rows), 'window': days, 'from': sent, 'to': end}