Postal v1.1: crash resume, A/B subjects, send-hours window, segments, suppression, unsub endpoint, watchdog, Proton SMTP
This commit is contained in:
444
engine.py
Normal file
444
engine.py
Normal file
@@ -0,0 +1,444 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Mailer Engine — Ollama-personalized email campaign sender.
|
||||
Runs on Windows (VPS). SQLite-backed. SES delivery. Ollama over Tailscale.
|
||||
"""
|
||||
import sqlite3, json, time, threading, queue, os, re, html, logging
|
||||
from datetime import datetime
|
||||
import requests
|
||||
|
||||
APP_DIR = os.path.dirname(os.path.abspath(__file__))
|
||||
DB_PATH = os.path.join(APP_DIR, "mailer.db")
|
||||
LOG_PATH = os.path.join(APP_DIR, "engine.log")
|
||||
|
||||
logging.basicConfig(filename=LOG_PATH, level=logging.INFO,
|
||||
format="%(asctime)s %(levelname)s %(message)s")
|
||||
log = logging.getLogger("mailer")
|
||||
|
||||
# ------------------------------------------------------------------ config
|
||||
def load_config():
|
||||
defaults = {
|
||||
"ollama_url": "http://100.103.34.0:11440", # nightmare tailnet via non-chunked shim
|
||||
"ollama_model": "mgraffam/gemma4-heretic:12b",
|
||||
"ses_region": "us-east-1",
|
||||
"send_mode": "smtp", # smtp (works now) | ses (needs AWS keys)
|
||||
"smtp_host": "smtp.protonmail.ch",
|
||||
"smtp_port": 587,
|
||||
"smtp_user": "makemoneys8@proton.me",
|
||||
"smtp_pass": "",
|
||||
"from_name": "drjones",
|
||||
"from_email": "",
|
||||
"reply_to": "",
|
||||
"rate_per_minute": 60, # SES default 14/s; keep low for reputation
|
||||
"batch_size": 50,
|
||||
"pause_between_batches_sec": 20,
|
||||
"personalize": True,
|
||||
"personalize_prompt": "Write a 2-3 sentence friendly, personal email greeting for {first_name}. They are a subscriber interested in {interest}. Warm, human, no hype, no exclamation marks. Sign as {from_name}.",
|
||||
"unsubscribe_footer": True,
|
||||
"send_hours": "9-18", # server-local hour window; outside = pause
|
||||
"ab_subjects": True, # gemma4 writes 3 subject variants
|
||||
"stopped": False,
|
||||
}
|
||||
c = sqlite3.connect(DB_PATH)
|
||||
c.execute("CREATE TABLE IF NOT EXISTS config (k TEXT PRIMARY KEY, v TEXT)")
|
||||
for k, v in defaults.items():
|
||||
c.execute("INSERT OR IGNORE INTO config (k,v) VALUES (?,?)", (k, json.dumps(v)))
|
||||
c.commit()
|
||||
out = {k: json.loads(v) for (k, v) in c.execute("SELECT k,v FROM config")}
|
||||
c.close()
|
||||
return out
|
||||
|
||||
def save_config(cfg):
|
||||
c = sqlite3.connect(DB_PATH)
|
||||
for k, v in cfg.items():
|
||||
c.execute("INSERT OR REPLACE INTO config (k,v) VALUES (?,?)", (k, json.dumps(v)))
|
||||
c.commit(); c.close()
|
||||
|
||||
# ------------------------------------------------------------------ db
|
||||
def db():
|
||||
c = sqlite3.connect(DB_PATH, check_same_thread=False)
|
||||
c.row_factory = sqlite3.Row
|
||||
return c
|
||||
|
||||
def init_db():
|
||||
c = db()
|
||||
c.executescript("""
|
||||
CREATE TABLE IF NOT EXISTS subscribers (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
email TEXT UNIQUE NOT NULL,
|
||||
first_name TEXT DEFAULT '',
|
||||
last_name TEXT DEFAULT '',
|
||||
interest TEXT DEFAULT '',
|
||||
status TEXT DEFAULT 'active', -- active|unsubscribed|bounced|complained
|
||||
source TEXT DEFAULT 'import',
|
||||
created_at TEXT DEFAULT (datetime('now'))
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS sends (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
subscriber_id INTEGER,
|
||||
campaign_id INTEGER,
|
||||
email TEXT,
|
||||
subject TEXT,
|
||||
body TEXT,
|
||||
status TEXT DEFAULT 'queued', -- queued|personalizing|sent|failed|skipped
|
||||
error TEXT DEFAULT '',
|
||||
message_id TEXT DEFAULT '',
|
||||
sent_at TEXT
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS campaigns (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name TEXT,
|
||||
subject TEXT,
|
||||
body TEXT, -- HTML with {{first_name}} etc
|
||||
created_at TEXT DEFAULT (datetime('now')),
|
||||
status TEXT DEFAULT 'draft', -- draft|ready|running|paused|done
|
||||
segment_interest TEXT DEFAULT ''
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS suppression (
|
||||
email TEXT PRIMARY KEY,
|
||||
reason TEXT DEFAULT 'manual',
|
||||
created_at TEXT DEFAULT (datetime('now'))
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS ab_variants (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
campaign_id INTEGER,
|
||||
subject TEXT,
|
||||
sent INTEGER DEFAULT 0,
|
||||
opened INTEGER DEFAULT 0
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS events (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
send_id INTEGER,
|
||||
type TEXT, -- open|click|bounce|complaint|unsubscribe
|
||||
ts TEXT DEFAULT (datetime('now'))
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_sends_status ON sends(status);
|
||||
CREATE INDEX IF NOT EXISTS idx_subs_status ON subscribers(status);
|
||||
""")
|
||||
c.commit(); c.close()
|
||||
|
||||
# ------------------------------------------------------------------ ollama
|
||||
class OllamaClient:
|
||||
def __init__(self, url, model, timeout=300):
|
||||
self.url = url.rstrip("/")
|
||||
self.model = model
|
||||
self.timeout = timeout
|
||||
self._fail_until = 0
|
||||
|
||||
def healthy(self):
|
||||
try:
|
||||
r = requests.get(f"{self.url}/api/tags", timeout=5)
|
||||
return r.status_code == 200
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
def generate(self, prompt):
|
||||
if time.time() < self._fail_until:
|
||||
return None
|
||||
try:
|
||||
r = requests.post(f"{self.url}/api/generate",
|
||||
json={"model": self.model, "prompt": prompt, "stream": False,
|
||||
"options": {"num_predict": 350, "temperature": 0.8}, "keep_alive": "15m"},
|
||||
timeout=self.timeout)
|
||||
r.raise_for_status()
|
||||
d = r.json()
|
||||
text = (d.get("response") or "").strip()
|
||||
if not text:
|
||||
text = (d.get("thinking") or "").strip()
|
||||
# strip any raw-CoT
|
||||
if "</think>" in text:
|
||||
text = text.split("</think>")[-1]
|
||||
return text.strip()
|
||||
except Exception as e:
|
||||
log.warning(f"ollama fail: {e}")
|
||||
self._fail_until = time.time() + 120 # backoff 2 min
|
||||
return None
|
||||
|
||||
def personalize_intro(ollama, sub, cfg):
|
||||
if not cfg.get("personalize") or not ollama:
|
||||
return ""
|
||||
prompt = ("Reply with ONLY the greeting text itself — no plan, no options, no bullets, no thinking out loud.\n"
|
||||
+ cfg["personalize_prompt"].format(
|
||||
first_name=sub["first_name"] or "there",
|
||||
interest=sub["interest"] or "our updates",
|
||||
from_name=cfg.get("from_name", "")))
|
||||
out = ollama.generate(prompt)
|
||||
if not out:
|
||||
return ""
|
||||
# sanitize: keep it short, no subject lines the model might add
|
||||
# gemma4-heretic wraps candidate greetings in "double quotes" inside planning bullets.
|
||||
quoted = re.findall(r'"([^"\n]{25,400})"', out)
|
||||
if quoted:
|
||||
text = quoted[-1].strip()
|
||||
else:
|
||||
keep = []
|
||||
for l in out.splitlines():
|
||||
t = l.strip()
|
||||
if not t:
|
||||
continue
|
||||
if t.startswith("*") or t.startswith("-"):
|
||||
continue
|
||||
low = t.lower()
|
||||
if low.startswith(("constraint", "subject:", "tone:", "content:", "format",
|
||||
"option", "plan", "context", "recipient", "interest",
|
||||
"sign-off", "sign:")):
|
||||
continue
|
||||
keep.append(t)
|
||||
text = " ".join(keep) if keep else out.strip()
|
||||
return text[:600]
|
||||
|
||||
# ------------------------------------------------------------------ template
|
||||
def render_template(body, sub, cfg, intro=""):
|
||||
d = sub["first_name"] or "there"
|
||||
out = body
|
||||
for k, v in {
|
||||
"{{first_name}}": d,
|
||||
"{{first_name_or_there}}": d,
|
||||
"{{last_name}}": sub["last_name"] or "",
|
||||
"{{email}}": sub["email"],
|
||||
"{{interest}}": sub["interest"] or "our updates",
|
||||
"{{from_name}}": cfg.get("from_name", ""),
|
||||
"{{intro}}": intro,
|
||||
}.items():
|
||||
out = out.replace(k, html.escape(v) if "{{intro}}" != k else v)
|
||||
# also {first_name} single-brace style
|
||||
out = re.sub(r"\{first_name\}", d, out)
|
||||
if cfg.get("unsubscribe_footer"):
|
||||
foot = ('<div style="margin-top:32px;padding-top:12px;border-top:1px solid #ddd;'
|
||||
'font-size:12px;color:#999">You get this because you subscribed. '
|
||||
'<a href="{{UNSUB}}">Unsubscribe</a></div>')
|
||||
out = out + foot
|
||||
return out
|
||||
|
||||
def render_subject(subject, sub):
|
||||
out = subject.replace("{{first_name}}", sub["first_name"] or "there")
|
||||
out = re.sub(r"\{first_name\}", sub["first_name"] or "there", out)
|
||||
return out
|
||||
|
||||
# ------------------------------------------------------------------ sender
|
||||
class Sender:
|
||||
"""SES via boto3; lazy import so the GUI can start without creds."""
|
||||
def __init__(self, cfg):
|
||||
self.cfg = cfg
|
||||
self._client = None
|
||||
|
||||
def client(self):
|
||||
if self._client is None:
|
||||
import boto3
|
||||
self._client = boto3.client("sesv2", region_name=self.cfg.get("ses_region", "us-east-1"))
|
||||
return self._client
|
||||
|
||||
def send(self, to_email, subject, html_body):
|
||||
if self.cfg.get("send_mode") == "smtp":
|
||||
return self._send_smtp(to_email, subject, html_body)
|
||||
c = self.client()
|
||||
resp = c.send_email(
|
||||
FromEmailAddress=f'{self.cfg.get("from_name","")} <{self.cfg["from_email"]}>' if self.cfg.get("from_name") else self.cfg["from_email"],
|
||||
Destination={"ToAddresses": [to_email]},
|
||||
Content={
|
||||
"Simple": {
|
||||
"Subject": {"Data": subject, "Charset": "UTF-8"},
|
||||
"Body": {"Html": {"Data": html_body, "Charset": "UTF-8"}},
|
||||
}
|
||||
},
|
||||
ReplyToAddresses=[self.cfg["reply_to"]] if self.cfg.get("reply_to") else [],
|
||||
)
|
||||
return resp["MessageId"]
|
||||
|
||||
def _send_smtp(self, to_email, subject, html_body):
|
||||
import smtplib
|
||||
from email.mime.multipart import MIMEMultipart
|
||||
from email.mime.text import MIMEText
|
||||
msg = MIMEMultipart("alternative")
|
||||
msg["Subject"] = subject
|
||||
msg["From"] = f'{self.cfg.get("from_name","")} <{self.cfg.get("smtp_user","")}>'
|
||||
msg["To"] = to_email
|
||||
if self.cfg.get("reply_to"):
|
||||
msg["Reply-To"] = self.cfg["reply_to"]
|
||||
msg.attach(MIMEText(html_body, "html", "utf-8"))
|
||||
with smtplib.SMTP(self.cfg.get("smtp_host", "smtp.mail.me.com"),
|
||||
int(self.cfg.get("smtp_port", 587)), timeout=30) as srv:
|
||||
srv.starttls()
|
||||
srv.login(self.cfg.get("smtp_user", ""), self.cfg.get("smtp_pass", ""))
|
||||
srv.send_message(msg)
|
||||
return f"smtp-{int(time.time())}"
|
||||
|
||||
# ------------------------------------------------------------------ engine
|
||||
class Engine:
|
||||
def __init__(self):
|
||||
init_db()
|
||||
self.state = {"running": False, "current": "", "sent": 0, "failed": 0,
|
||||
"skipped": 0, "total": 0, "started_at": ""}
|
||||
self.thread = None
|
||||
|
||||
def status(self):
|
||||
c = db()
|
||||
counts = {r["status"]: r["n"] for r in c.execute(
|
||||
"SELECT status, COUNT(*) n FROM sends GROUP BY status")}
|
||||
camp = c.execute("SELECT * FROM campaigns ORDER BY id DESC LIMIT 1").fetchone()
|
||||
c.close()
|
||||
return {**self.state, "counts": counts,
|
||||
"campaign": dict(camp) if camp else None}
|
||||
|
||||
def run_campaign(self, campaign_id):
|
||||
if self.state["running"]:
|
||||
return {"ok": False, "error": "already running"}
|
||||
self.thread = threading.Thread(target=self._run, args=(campaign_id,), daemon=True)
|
||||
self.thread.start()
|
||||
return {"ok": True}
|
||||
|
||||
def _in_send_hours(self, cfg):
|
||||
try:
|
||||
lo, hi = cfg.get("send_hours", "9-18").split("-")
|
||||
h = datetime.now().hour
|
||||
return int(lo) <= h < int(hi)
|
||||
except Exception:
|
||||
return True
|
||||
|
||||
def _run(self, campaign_id):
|
||||
cfg = load_config()
|
||||
c = db()
|
||||
camp = c.execute("SELECT * FROM campaigns WHERE id=?", (campaign_id,)).fetchone()
|
||||
if not camp:
|
||||
c.close(); return
|
||||
# queue all active subscribers not yet sent for this campaign
|
||||
seg = (camp["segment_interest"] or "").strip() if "segment_interest" in camp.keys() else ""
|
||||
if seg:
|
||||
rows = c.execute("""
|
||||
SELECT s.* FROM subscribers s
|
||||
WHERE s.status='active'
|
||||
AND s.interest LIKE ?
|
||||
AND s.id NOT IN (SELECT subscriber_id FROM sends WHERE campaign_id=? AND status IN ('sent','skipped'))
|
||||
ORDER BY s.id
|
||||
""", (f"%{seg}%", campaign_id)).fetchall()
|
||||
else:
|
||||
rows = c.execute("""
|
||||
SELECT s.* FROM subscribers s
|
||||
WHERE s.status='active'
|
||||
AND s.id NOT IN (SELECT subscriber_id FROM sends WHERE campaign_id=? AND status IN ('sent','skipped'))
|
||||
ORDER BY s.id
|
||||
""", (campaign_id,)).fetchall()
|
||||
for r in rows:
|
||||
c.execute("INSERT OR IGNORE INTO sends (subscriber_id, campaign_id, email, status) VALUES (?,?,?,'queued')",
|
||||
(r["id"], campaign_id, r["email"]))
|
||||
c.commit()
|
||||
self.state.update({"running": True, "sent": 0, "failed": 0, "skipped": 0,
|
||||
"total": len(rows), "started_at": datetime.now().isoformat(timespec='seconds')})
|
||||
log.info(f"campaign {campaign_id} start: {len(rows)} recipients")
|
||||
|
||||
ollama = OllamaClient(cfg["ollama_url"], cfg["ollama_model"]) if cfg.get("personalize") else None
|
||||
if ollama:
|
||||
self.state["current"] = f"checking Ollama at {cfg['ollama_url']}..."
|
||||
if ollama.healthy():
|
||||
log.info("ollama healthy")
|
||||
else:
|
||||
log.warning("ollama NOT reachable — sending without personalization")
|
||||
self.state["current"] = "Ollama unreachable — sending generic (will retry it periodically)"
|
||||
|
||||
# A/B subject variants (feature 6): gemma4 writes 2 alternates; variant 0 = campaign subject
|
||||
variants = [camp["subject"]]
|
||||
if cfg.get("ab_subjects") and ollama and ollama.healthy():
|
||||
for i in range(2):
|
||||
alt = ollama.generate(
|
||||
f"Rewrite this email subject line differently — same meaning, fresh angle, under 60 chars, no quotes. Reply with ONLY the subject.\nSubject: {camp['subject']}")
|
||||
if alt:
|
||||
alt = alt.strip().strip('"').splitlines()[0][:80]
|
||||
if alt and alt.lower() not in [v.lower() for v in variants]:
|
||||
variants.append(alt)
|
||||
if len(variants) > 1:
|
||||
c.execute("DELETE FROM ab_variants WHERE campaign_id=?", (campaign_id,))
|
||||
for v in variants:
|
||||
c.execute("INSERT INTO ab_variants (campaign_id, subject) VALUES (?,?)", (campaign_id, v))
|
||||
c.commit()
|
||||
log.info(f"AB variants: {variants}")
|
||||
|
||||
sender = Sender(cfg)
|
||||
batch_n = 0
|
||||
try:
|
||||
sends = c.execute("SELECT * FROM sends WHERE campaign_id=? AND status='queued' ORDER BY id",
|
||||
(campaign_id,)).fetchall()
|
||||
for snd in sends:
|
||||
if cfg.get("stopped") or not self.state["running"]:
|
||||
log.info("stopped by operator")
|
||||
break
|
||||
sub = c.execute("SELECT * FROM subscribers WHERE id=?", (snd["subscriber_id"],)).fetchone()
|
||||
if not sub or sub["status"] != "active":
|
||||
c.execute("UPDATE sends SET status='skipped' WHERE id=?", (snd["id"],)); c.commit()
|
||||
self.state["skipped"] += 1
|
||||
continue
|
||||
if c.execute("SELECT 1 FROM suppression WHERE email=?", (sub["email"],)).fetchone():
|
||||
c.execute("UPDATE sends SET status='skipped', error='suppressed' WHERE id=?", (snd["id"],)); c.commit()
|
||||
self.state["skipped"] += 1
|
||||
continue
|
||||
if not self._in_send_hours(cfg):
|
||||
# outside sending window: pause loop (re-check every 5 min)
|
||||
self.state["current"] = f"outside send hours ({cfg.get('send_hours')}) — sleeping"
|
||||
c.execute("UPDATE sends SET status='queued' WHERE id=?", (snd["id"],)); c.commit()
|
||||
time.sleep(300)
|
||||
continue
|
||||
self.state["current"] = sub["email"]
|
||||
# personalize
|
||||
intro = ""
|
||||
if ollama and cfg.get("personalize"):
|
||||
c.execute("UPDATE sends SET status='personalizing' WHERE id=?", (snd["id"],)); c.commit()
|
||||
intro = personalize_intro(ollama, sub, cfg)
|
||||
variant_i = self.state["sent"] % len(variants)
|
||||
subject = render_subject(variants[variant_i], sub)
|
||||
body = render_template(camp["body"], sub, cfg, intro)
|
||||
# send
|
||||
try:
|
||||
mid = sender.send(sub["email"], subject, body)
|
||||
c.execute("UPDATE sends SET status='sent', subject=?, body=?, message_id=?, sent_at=datetime('now') WHERE id=?",
|
||||
(subject, body, mid, snd["id"]))
|
||||
if len(variants) > 1:
|
||||
c.execute("UPDATE ab_variants SET sent=sent+1 WHERE campaign_id=? AND subject=?",
|
||||
(campaign_id, variants[variant_i]))
|
||||
self.state["sent"] += 1
|
||||
except Exception as e:
|
||||
err = str(e)[:300]
|
||||
log.error(f"send fail {sub['email']}: {err}")
|
||||
c.execute("UPDATE sends SET status='failed', error=? WHERE id=?", (err, snd["id"]))
|
||||
self.state["failed"] += 1
|
||||
if "Throttling" in err:
|
||||
time.sleep(30)
|
||||
c.commit()
|
||||
batch_n += 1
|
||||
time.sleep(60.0 / max(1, cfg.get("rate_per_minute", 60)))
|
||||
if batch_n >= cfg.get("batch_size", 50):
|
||||
batch_n = 0
|
||||
self.state["current"] = f"batch done — pausing {cfg.get('pause_between_batches_sec',20)}s"
|
||||
time.sleep(cfg.get("pause_between_batches_sec", 20))
|
||||
finally:
|
||||
self.state["running"] = False
|
||||
self.state["current"] = "idle"
|
||||
c.execute("UPDATE campaigns SET status='done' WHERE id=?", (campaign_id,))
|
||||
c.commit(); c.close()
|
||||
log.info(f"campaign {campaign_id} done: sent={self.state['sent']} failed={self.state['failed']}")
|
||||
|
||||
ENGINE = None
|
||||
def get_engine():
|
||||
global ENGINE
|
||||
if ENGINE is None:
|
||||
ENGINE = Engine()
|
||||
ENGINE.auto_resume()
|
||||
return ENGINE
|
||||
|
||||
def auto_resume_thread():
|
||||
"""Called from app startup: find campaign marked running/ready with queued sends and restart."""
|
||||
import time as _t
|
||||
_t.sleep(15) # let the web app bind first
|
||||
try:
|
||||
eng = get_engine()
|
||||
c = db()
|
||||
row = c.execute("SELECT id FROM campaigns WHERE status IN ('running','ready') ORDER BY id DESC LIMIT 1").fetchone()
|
||||
n = c.execute("SELECT COUNT(*) n FROM sends WHERE status='queued'").fetchone()["n"]
|
||||
c.close()
|
||||
if row and n and not eng.state["running"]:
|
||||
log.info(f"auto-resume campaign {row['id']} ({n} queued)")
|
||||
eng.run_campaign(row["id"])
|
||||
except Exception as e:
|
||||
log.error(f"auto-resume failed: {e}")
|
||||
|
||||
import threading as _threading
|
||||
def start_auto_resume():
|
||||
_threading.Thread(target=auto_resume_thread, daemon=True).start()
|
||||
Reference in New Issue
Block a user