"""Auto-extracted from models.py — do not edit manually."""
import json
import logging
import secrets
import sqlite3
import uuid
from datetime import UTC, datetime, timezone
import bcrypt
logger = logging.getLogger(__name__)
from app.crypto import (
decrypt_entity,
decrypt_user_value,
encrypt_entity,
encrypt_user_submission,
encrypt_user_value,
encrypt_value,
hash_value,
is_encrypted,
key_is_configured,
try_decrypt,
try_decrypt_entity,
try_decrypt_user_submission,
try_decrypt_user_value,
)
from app.db import DB_PATH, get_db
# ─── Column-name allowlist helpers ──────────────────────────────────────
_CAMPAIGN_RECIPIENTS_ALLOWED_COLUMNS = {
"status",
"email",
"metadata",
"sent_at",
"opened_at",
"clicked_at",
"message_id",
"error",
"bounce_type",
"bounce_reason",
"tracking_token",
"campaign_id",
"created_at",
}
def _validate_columns(columns: list, allowed: set) -> list:
"""Validate column names against an allowlist. Returns list of 'col = ?' clauses."""
valid = []
for col in columns:
if col in allowed:
valid.append(f"{col} = ?")
else:
logger.warning(f"Blocked column in UPDATE: {col} (not in allowlist)")
return valid
def get_custom_templates(user_id):
"""Get all custom templates for a user."""
conn = get_db()
try:
rows = conn.execute(
"SELECT * FROM custom_templates WHERE user_id = ? ORDER BY name",
(user_id,),
).fetchall()
return [dict(row) for row in rows]
finally:
conn.close()
def get_custom_template(user_id, name):
"""Get a specific custom template by name."""
conn = get_db()
try:
row = conn.execute(
"SELECT * FROM custom_templates WHERE user_id = ? AND name = ?",
(user_id, name),
).fetchone()
return dict(row) if row else None
finally:
conn.close()
def create_custom_template(user_id, name, content, description=None):
"""Create a new custom template."""
conn = get_db()
try:
cursor = conn.execute(
"INSERT INTO custom_templates (user_id, name, content, description) VALUES (?, ?, ?, ?)",
(user_id, name, content, description),
)
conn.commit()
return get_custom_template(user_id, name)
except sqlite3.IntegrityError:
return get_custom_template(user_id, name) # Already exists
finally:
conn.close()
def update_custom_template(user_id, name, content=None, description=None):
"""Update a custom template."""
conn = get_db()
try:
existing = get_custom_template(user_id, name)
if not existing:
return None
if content is not None:
conn.execute(
"UPDATE custom_templates SET content = ?, updated_at = CURRENT_TIMESTAMP WHERE user_id = ? AND name = ?",
(content, user_id, name),
)
if description is not None:
conn.execute(
"UPDATE custom_templates SET description = ?, updated_at = CURRENT_TIMESTAMP WHERE user_id = ? AND name = ?",
(description, user_id, name),
)
conn.commit()
return get_custom_template(user_id, name)
finally:
conn.close()
def delete_custom_template(user_id, name):
"""Delete a custom template."""
from app.emails.renderer import _get_renderer
conn = get_db()
try:
existing = get_custom_template(user_id, name)
if not existing:
return False
conn.execute(
"DELETE FROM custom_templates WHERE user_id = ? AND name = ?",
(user_id, name),
)
conn.commit()
# Also delete from disk
_get_renderer().delete_custom_template(name)
return True
finally:
conn.close()
def create_campaign(user_id, name, subject, body, from_name=None, from_email=None, template_id=None, scheduled_at=None):
"""Create a new email campaign. Returns campaign dict."""
conn = None
try:
conn = get_db()
cur = conn.execute(
"""INSERT INTO campaigns (user_id, name, subject, body, from_name, from_email, template_id, scheduled_at, status)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(
user_id,
name,
subject,
body,
from_name,
from_email,
template_id,
scheduled_at,
"scheduled" if scheduled_at else "draft",
),
)
conn.commit()
return get_campaign(cur.lastrowid, user_id)
finally:
if conn:
conn.close()
def get_campaign(campaign_id, user_id):
"""Get a single campaign. Returns dict or None."""
conn = None
try:
conn = get_db()
row = conn.execute(
"SELECT * FROM campaigns WHERE id = ? AND user_id = ?",
(campaign_id, user_id),
).fetchone()
if not row:
return None
camp = dict(row)
# Add stats
stats = conn.execute(
"""SELECT status, COUNT(*) as cnt FROM campaign_recipients
WHERE campaign_id = ? GROUP BY status""",
(campaign_id,),
).fetchall()
camp["stats"] = {s["status"]: s["cnt"] for s in stats}
camp["stats"]["total"] = sum(s["cnt"] for s in stats)
return camp
finally:
if conn:
conn.close()
def get_campaign_by_id(campaign_id):
"""Get campaign by ID without user check (for internal worker use). Returns dict or None."""
conn = None
try:
conn = get_db()
row = conn.execute(
"SELECT * FROM campaigns WHERE id = ?",
(campaign_id,),
).fetchone()
if not row:
return None
camp = dict(row)
stats = conn.execute(
"""SELECT status, COUNT(*) as cnt FROM campaign_recipients
WHERE campaign_id = ? GROUP BY status""",
(campaign_id,),
).fetchall()
camp["stats"] = {s["status"]: s["cnt"] for s in stats}
camp["stats"]["total"] = sum(s["cnt"] for s in stats)
return camp
finally:
if conn:
conn.close()
def list_campaigns(user_id, status=None, limit=50, offset=0):
"""List campaigns for a user. Returns list of dicts."""
conn = None
try:
conn = get_db()
query = "SELECT * FROM campaigns WHERE user_id = ?"
params = [user_id]
if status:
query += " AND status = ?"
params.append(status)
query += " ORDER BY created_at DESC LIMIT ? OFFSET ?"
params += [limit, offset]
rows = conn.execute(query, params).fetchall()
campaigns = [dict(r) for r in rows]
for camp in campaigns:
stats = conn.execute(
"""SELECT status, COUNT(*) as cnt FROM campaign_recipients
WHERE campaign_id = ? GROUP BY status""",
(camp["id"],),
).fetchall()
camp["stats"] = {s["status"]: s["cnt"] for s in stats}
camp["stats"]["total"] = sum(s["cnt"] for s in stats)
return campaigns
finally:
if conn:
conn.close()
def add_campaign_recipients(campaign_id, user_id, emails):
"""Add recipients to a campaign. Returns dict with counts."""
conn = None
try:
conn = get_db()
# Verify campaign belongs to user
camp = conn.execute(
"SELECT id FROM campaigns WHERE id = ? AND user_id = ?",
(campaign_id, user_id),
).fetchone()
if not camp:
return {"error": "campaign not found"}
added = 0
for email in emails:
token = secrets.token_urlsafe(32)
conn.execute(
"INSERT INTO campaign_recipients (campaign_id, email, status, tracking_token) VALUES (?, ?, 'pending', ?)",
(campaign_id, email, token),
)
added += 1
conn.commit()
return {"added": added, "campaign_id": campaign_id}
finally:
if conn:
conn.close()
def update_campaign_status(campaign_id, user_id, status):
"""Update campaign status. Returns updated campaign dict."""
conn = None
try:
conn = get_db()
conn.execute(
"UPDATE campaigns SET status = ?, updated_at = CURRENT_TIMESTAMP WHERE id = ? AND user_id = ?",
(status, campaign_id, user_id),
)
conn.commit()
return get_campaign(campaign_id, user_id)
finally:
if conn:
conn.close()
def delete_campaign(campaign_id, user_id):
"""Delete a campaign and its recipients. Returns True if deleted."""
conn = None
try:
conn = get_db()
camp = conn.execute(
"SELECT id FROM campaigns WHERE id = ? AND user_id = ?",
(campaign_id, user_id),
).fetchone()
if not camp:
return False
conn.execute("DELETE FROM campaign_recipients WHERE campaign_id = ?", (campaign_id,))
conn.execute("DELETE FROM campaigns WHERE id = ? AND user_id = ?", (campaign_id, user_id))
conn.commit()
return True
finally:
if conn:
conn.close()
def get_campaign_recipients(campaign_id, user_id, status=None, limit=100, offset=0):
"""Get recipients for a campaign. Returns list of dicts."""
conn = None
try:
conn = get_db()
# Verify campaign belongs to user
camp = conn.execute(
"SELECT id FROM campaigns WHERE id = ? AND user_id = ?",
(campaign_id, user_id),
).fetchone()
if not camp:
return []
query = """SELECT cr.*, s.name FROM campaign_recipients cr
JOIN campaigns s ON cr.campaign_id = s.id
WHERE cr.campaign_id = ?"""
params = [campaign_id]
if status:
query += " AND cr.status = ?"
params.append(status)
query += " ORDER BY cr.created_at DESC LIMIT ? OFFSET ?"
params += [limit, offset]
rows = conn.execute(query, params).fetchall()
return [dict(r) for r in rows]
finally:
if conn:
conn.close()
def get_pending_recipients(campaign_id):
"""Get pending recipients for a campaign (for worker). Returns list of dicts."""
conn = None
try:
conn = get_db()
rows = conn.execute(
"""SELECT cr.*, c.subject, c.body, c.from_name, c.from_email, c.template_id,
c.user_id
FROM campaign_recipients cr
JOIN campaigns c ON cr.campaign_id = c.id
WHERE cr.campaign_id = ? AND cr.status = 'pending'""",
(campaign_id,),
).fetchall()
return [dict(r) for r in rows]
finally:
if conn:
conn.close()
def update_recipient_status(campaign_id, recipient_id, status, message_id=None, error=None):
"""Update recipient status. Returns True if updated."""
conn = None
try:
conn = get_db()
fields = ["status = ?"]
params = [status]
if message_id:
fields.append("message_id = ?")
params.append(message_id)
if error:
fields.append("error = ?")
params.append(error)
if status == "sent":
fields.append("sent_at = CURRENT_TIMESTAMP")
params += [recipient_id, campaign_id]
conn.execute(
f"UPDATE campaign_recipients SET {', '.join(fields)} WHERE id = ? AND campaign_id = ?",
params,
)
conn.commit()
return True
finally:
if conn:
conn.close()
def get_recipient_by_message_id(message_id):
"""Find a campaign recipient by Resend message_id.
Searches across all campaigns since we don't know which campaign this
message belongs to from the webhook event alone.
Returns recipient dict or None.
"""
conn = None
try:
conn = get_db()
row = conn.execute(
"SELECT id, campaign_id, user_id, email FROM campaign_recipients "
"WHERE message_id = ? ORDER BY created_at DESC LIMIT 1",
(message_id,),
).fetchone()
return dict(row) if row else None
finally:
if conn:
conn.close()
def get_email_settings(user_id):
"""Get email settings for a user. Returns dict or None."""
conn = None
try:
conn = get_db()
row = conn.execute(
"SELECT * FROM email_settings WHERE user_id = ?",
(user_id,),
).fetchone()
return dict(row) if row else None
finally:
if conn:
conn.close()
# ─── Email Tracking ──────────────────────────────────────────────────────────
def get_recipient_by_token(tracking_token):
"""Look up a campaign recipient by tracking token. Returns dict or None."""
conn = None
try:
conn = get_db()
row = conn.execute(
"SELECT cr.*, c.user_id, c.subject, c.from_email FROM campaign_recipients cr "
"JOIN campaigns c ON cr.campaign_id = c.id "
"WHERE cr.tracking_token = ?",
(tracking_token,),
).fetchone()
return dict(row) if row else None
finally:
if conn:
conn.close()
def record_open(recipient_id):
"""Record an email open event."""
conn = None
try:
conn = get_db()
conn.execute(
"UPDATE campaign_recipients SET opened_at = COALESCE(opened_at, ?) WHERE id = ?",
(datetime.now(UTC).isoformat(), recipient_id),
)
conn.commit()
finally:
if conn:
conn.close()
def record_click(recipient_id):
"""Record a link click event."""
conn = None
try:
conn = get_db()
conn.execute(
"UPDATE campaign_recipients SET clicked_at = COALESCE(clicked_at, ?) WHERE id = ?",
(datetime.now(UTC).isoformat(), recipient_id),
)
conn.commit()
finally:
if conn:
conn.close()
def record_bounce(recipient_id, bounce_type, bounce_reason=None):
"""Record a bounce event."""
conn = None
try:
conn = get_db()
conn.execute(
"UPDATE campaign_recipients SET status = 'bounced', bounce_type = ?, bounce_reason = ? WHERE id = ?",
(bounce_type, bounce_reason, recipient_id),
)
conn.commit()
finally:
if conn:
conn.close()
def get_campaign_recipients_by_email(email):
"""Find all recipients matching an email address (for bounce lookups)."""
conn = None
try:
conn = get_db()
rows = conn.execute(
"SELECT cr.* FROM campaign_recipients cr WHERE LOWER(cr.email) = ?",
(email.lower(),),
).fetchall()
return [dict(r) for r in rows]
finally:
if conn:
conn.close()
def get_campaign_analytics(campaign_id, user_id):
"""Get open/click/bounce analytics for a campaign."""
conn = None
try:
conn = get_db()
# Verify ownership
owner = conn.execute(
"SELECT user_id FROM campaigns WHERE id = ?",
(campaign_id,),
).fetchone()
if not owner or owner["user_id"] != user_id:
return None
row = conn.execute(
"SELECT "
" COUNT(*) as total_recipients, "
" SUM(CASE WHEN status = 'sent' THEN 1 ELSE 0 END) as sent, "
" SUM(CASE WHEN status = 'failed' THEN 1 ELSE 0 END) as failed, "
" SUM(CASE WHEN status = 'bounced' THEN 1 ELSE 0 END) as bounced, "
" SUM(CASE WHEN opened_at IS NOT NULL THEN 1 ELSE 0 END) as unique_opens, "
" SUM(CASE WHEN clicked_at IS NOT NULL THEN 1 ELSE 0 END) as unique_clicks "
"FROM campaign_recipients "
"WHERE campaign_id = ?",
(campaign_id,),
).fetchone()
if not row:
return None
result = dict(row)
total = result.get("total_recipients", 0) or 0
sent = result.get("sent", 0) or 0
opens = result.get("unique_opens", 0) or 0
clicks = result.get("unique_clicks", 0) or 0
bounces = result.get("bounced", 0) or 0
result["open_rate"] = round(opens / total, 2) if total else 0.0
result["click_rate"] = round(clicks / total, 2) if total else 0.0
result["bounce_rate"] = round(bounces / total, 2) if total else 0.0
# Aliases for consistency
result["unique_bounces"] = bounces
return result
finally:
if conn:
conn.close()
# ─── Form Notification Email Tracking ───────────────────────────────────────
def log_notification_email(user_id, submission_id, to_email, tracking_token):
"""Log a form notification email to email_send_log with tracking token.
Returns the log entry ID.
"""
conn = None
try:
conn = get_db()
result = conn.execute(
"INSERT INTO email_send_log (user_id, submission_id, tracking_token, count) VALUES (?, ?, ?, 1)",
(user_id, submission_id, tracking_token),
)
conn.commit()
return result.lastrowid
finally:
if conn:
conn.close()
def get_notification_email_by_token(tracking_token):
"""Look up a form notification email by tracking token. Returns dict or None."""
conn = None
try:
conn = get_db()
row = conn.execute(
"SELECT esl.* FROM email_send_log esl WHERE esl.tracking_token = ? AND esl.submission_id IS NOT NULL",
(tracking_token,),
).fetchone()
return dict(row) if row else None
finally:
if conn:
conn.close()
def record_notification_open(log_id):
"""Record an open for a form notification email."""
conn = None
try:
conn = get_db()
conn.execute(
"UPDATE email_send_log SET sent_at = COALESCE(sent_at, ?) WHERE id = ?",
(datetime.now(UTC).isoformat(), log_id),
)
conn.commit()
finally:
if conn:
conn.close()
def record_notification_click(log_id):
"""Record a click for a form notification email."""
conn = None
try:
conn = get_db()
conn.execute(
"UPDATE email_send_log SET count = count + 1 WHERE id = ?",
(log_id,),
)
conn.commit()
finally:
if conn:
conn.close()
def get_user_email_analytics(user_id):
"""Get email analytics for a user's form notifications.
Returns dict with sent, opened, clicked counts.
"""
conn = None
try:
conn = get_db()
row = conn.execute(
"SELECT COUNT(*) as sent, "
"SUM(CASE WHEN submitted_at IS NOT NULL THEN 1 ELSE 0 END) as opened, "
"SUM(CASE WHEN count > 1 THEN 1 ELSE 0 END) as clicked "
"FROM email_send_log "
"WHERE user_id = ? AND submission_id IS NOT NULL",
(user_id,),
).fetchone()
return dict(row) if row else {"sent": 0, "opened": 0, "clicked": 0}
finally:
if conn:
conn.close()
def upsert_email_settings(
user_id,
from_name=None,
from_email=None,
reply_to=None,
template_id=None,
sender_domain=None,
domain_verified=None,
default_template=None,
):
"""Create or update email settings. Returns settings dict."""
conn = None
try:
conn = get_db()
existing = conn.execute(
"SELECT id FROM email_settings WHERE user_id = ?",
(user_id,),
).fetchone()
if existing:
conn.execute(
"""UPDATE email_settings SET from_name = ?, from_email = ?, reply_to = ?,
template_id = ?, sender_domain = ?, domain_verified = ?, default_template = ?,
updated_at = CURRENT_TIMESTAMP
WHERE user_id = ?""",
(
from_name,
from_email,
reply_to,
template_id,
sender_domain,
domain_verified,
default_template,
user_id,
),
)
else:
conn.execute(
"""INSERT INTO email_settings (user_id, from_name, from_email, reply_to, template_id, sender_domain, domain_verified, default_template)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)""",
(
user_id,
from_name,
from_email,
reply_to,
template_id,
sender_domain,
domain_verified,
default_template,
),
)
conn.commit()
return get_email_settings(user_id)
finally:
if conn:
conn.close()
# ─── Phase 4: Reminder Campaigns, A/B Testing, Rate Limiting ──────────────────────
# ── Reminder Campaigns ─────────────────────────────────────────────────────────────
def create_reminder(campaign_id, user_id, delay_hours, subject, body, template_id=None):
"""Create a reminder for a campaign. Returns reminder dict."""
conn = None
try:
conn = get_db()
# Verify campaign ownership
campaign = conn.execute(
"SELECT id FROM campaigns WHERE id = ? AND user_id = ?",
(campaign_id, user_id),
).fetchone()
if not campaign:
return None
cur = conn.execute(
"""INSERT INTO campaign_reminders (campaign_id, delay_hours, subject, body, template_id, status)
VALUES (?, ?, ?, ?, ?, 'scheduled')""",
(campaign_id, delay_hours, subject, body, template_id),
)
conn.commit()
return get_reminder(cur.lastrowid, user_id)
finally:
if conn:
conn.close()
def get_reminder(reminder_id, user_id):
"""Get a single reminder. Returns dict or None."""
conn = None
try:
conn = get_db()
row = conn.execute(
"""SELECT r.* FROM campaign_reminders r
JOIN campaigns c ON r.campaign_id = c.id
WHERE r.id = ? AND c.user_id = ?""",
(reminder_id, user_id),
).fetchone()
return dict(row) if row else None
finally:
if conn:
conn.close()
def get_reminders(campaign_id, user_id):
"""Get all reminders for a campaign. Returns list of dicts."""
conn = None
try:
conn = get_db()
rows = conn.execute(
"""SELECT r.* FROM campaign_reminders r
JOIN campaigns c ON r.campaign_id = c.id
WHERE r.campaign_id = ? AND c.user_id = ?
ORDER BY r.delay_hours""",
(campaign_id, user_id),
).fetchall()
return [dict(r) for r in rows]
finally:
if conn:
conn.close()
def cancel_reminder(reminder_id, user_id):
"""Cancel a scheduled reminder. Returns True if cancelled."""
conn = None
try:
conn = get_db()
result = conn.execute(
"""UPDATE campaign_reminders SET status = 'cancelled'
WHERE id = ? AND status = 'scheduled'
AND campaign_id IN (SELECT id FROM campaigns WHERE user_id = ?)""",
(reminder_id, user_id),
)
conn.commit()
return result.rowcount > 0
finally:
if conn:
conn.close()
def schedule_reminders_for_campaign(campaign_id, user_id):
"""Schedule all pending reminders after a campaign is sent.
Sets scheduled_at based on campaign's sent_at + delay_hours.
Returns list of scheduled reminders.
"""
conn = None
try:
conn = get_db()
# Get campaign's sent_at from first sent recipient
sent_at = conn.execute(
"""SELECT MIN(sent_at) as sent_at FROM campaign_recipients
WHERE campaign_id = ? AND status = 'sent'""",
(campaign_id,),
).fetchone()["sent_at"]
if not sent_at:
return []
# Also get campaign's user_id
campaign = conn.execute(
"SELECT user_id FROM campaigns WHERE id = ?",
(campaign_id,),
).fetchone()
if not campaign:
return []
reminders = conn.execute(
"SELECT * FROM campaign_reminders WHERE campaign_id = ? AND status = 'scheduled'",
(campaign_id,),
).fetchall()
scheduled = []
for r in reminders:
# Calculate scheduled_at
from datetime import datetime, timedelta
base = datetime.fromisoformat(str(sent_at))
scheduled_time = base + timedelta(hours=r["delay_hours"])
conn.execute(
"UPDATE campaign_reminders SET scheduled_at = ? WHERE id = ?",
(scheduled_time.isoformat(), r["id"]),
)
r = dict(r)
r["scheduled_at"] = scheduled_time.isoformat()
scheduled.append(r)
conn.commit()
return scheduled
finally:
if conn:
conn.close()
def get_pending_reminders():
"""Get all reminders that are scheduled and due. Returns list of dicts."""
from datetime import datetime
conn = None
try:
conn = get_db()
now = datetime.now(UTC).isoformat()
rows = conn.execute(
"""SELECT * FROM campaign_reminders
WHERE status = 'scheduled'
AND scheduled_at IS NOT NULL
AND scheduled_at <= ?""",
(now,),
).fetchall()
return [dict(r) for r in rows]
finally:
if conn:
conn.close()
def get_reminder_eligible_recipients(reminder_id):
"""Get recipients eligible for this reminder (haven't opened/clicked, haven't gotten reminder).
Returns list of recipient dicts.
"""
conn = None
try:
conn = get_db()
campaign_id = conn.execute(
"SELECT campaign_id FROM campaign_reminders WHERE id = ?",
(reminder_id,),
).fetchone()["campaign_id"]
rows = conn.execute(
"""SELECT * FROM campaign_recipients
WHERE campaign_id = ?
AND status = 'sent'
AND (opened_at IS NULL OR opened_at = '')
AND (clicked_at IS NULL OR clicked_at = '')
AND reminder_sent = 0""",
(campaign_id,),
).fetchall()
return [dict(r) for r in rows]
finally:
if conn:
conn.close()
def mark_reminder_sent(reminder_id, recipient_count):
"""Mark a reminder as sent and record recipient count."""
conn = None
try:
conn = get_db()
conn.execute(
"""UPDATE campaign_reminders
SET status = 'sent', sent_at = CURRENT_TIMESTAMP, recipient_count = ?
WHERE id = ?""",
(recipient_count, reminder_id),
)
conn.commit()
finally:
if conn:
conn.close()
# ── Rate Limiting ──────────────────────────────────────────────────────────────────
def check_send_quota(user_id):
"""Check if user has send quota remaining. Returns dict with quota info."""
conn = None
try:
conn = get_db()
settings = conn.execute(
"SELECT * FROM email_settings WHERE user_id = ?",
(user_id,),
).fetchone()
if not settings:
# No settings — use defaults
return {
"max_per_hour": 1000,
"max_per_day": 10000,
"used_this_hour": 0,
"used_this_day": 0,
"hour_reset_at": None,
"day_reset_at": None,
"can_send": True,
}
settings = dict(settings)
max_hour = settings.get("max_sends_per_hour") or 1000
max_day = settings.get("max_sends_per_day") or 10000
used_hour = settings.get("sends_this_hour") or 0
used_day = settings.get("sends_this_day") or 0
return {
"max_per_hour": max_hour,
"max_per_day": max_day,
"used_this_hour": used_hour,
"used_this_day": used_day,
"hour_reset_at": settings.get("hour_reset_at"),
"day_reset_at": settings.get("day_reset_at"),
"can_send": used_hour < max_hour and used_day < max_day,
}
finally:
if conn:
conn.close()
def record_send(user_id, count=1, campaign_id=None):
"""Record email sends for rate limiting. Returns True if send allowed."""
from datetime import datetime, timedelta
conn = None
try:
conn = get_db()
now = datetime.now(UTC)
settings = conn.execute(
"SELECT * FROM email_settings WHERE user_id = ?",
(user_id,),
).fetchone()
if not settings:
# Create default settings
conn.execute(
"""INSERT INTO email_settings (user_id, max_sends_per_hour, max_sends_per_day,
max_burst_per_campaign, sends_this_hour, sends_this_day,
hour_reset_at, day_reset_at)
VALUES (?, 1000, 10000, 100, 0, 0, ?, ?)""",
(user_id, (now + timedelta(hours=1)).isoformat(), (now + timedelta(days=1)).isoformat()),
)
settings = conn.execute(
"SELECT * FROM email_settings WHERE user_id = ?",
(user_id,),
).fetchone()
settings = dict(settings)
max_hour = settings.get("max_sends_per_hour") or 1000
max_day = settings.get("max_sends_per_day") or 10000
# Reset counters if time window has passed
used_hour = 0
used_day = 0
hour_reset = now + timedelta(hours=1)
day_reset = now + timedelta(days=1)
hour_reset_at = settings.get("hour_reset_at")
day_reset_at = settings.get("day_reset_at")
if hour_reset_at:
try:
hr = datetime.fromisoformat(str(hour_reset_at))
if now < hr:
used_hour = settings.get("sends_this_hour") or 0
hour_reset = hr
except (ValueError, TypeError):
pass
if day_reset_at:
try:
dr = datetime.fromisoformat(str(day_reset_at))
if now < dr:
used_day = settings.get("sends_this_day") or 0
day_reset = dr
except (ValueError, TypeError):
pass
# Check quota
if used_hour + count > max_hour:
return False
if used_day + count > max_day:
return False
# Update counters
conn.execute(
"""UPDATE email_settings
SET sends_this_hour = ?, sends_this_day = ?,
hour_reset_at = ?, day_reset_at = ?
WHERE user_id = ?""",
(used_hour + count, used_day + count, hour_reset.isoformat(), day_reset.isoformat(), user_id),
)
# Log the send
conn.execute(
"""INSERT INTO email_send_log (user_id, count, campaign_id)
VALUES (?, ?, ?)""",
(user_id, count, campaign_id),
)
conn.commit()
return True
finally:
if conn:
conn.close()
def cleanup_old_send_log():
"""Remove send log entries older than 7 days."""
conn = None
try:
conn = get_db()
conn.execute(
"""DELETE FROM email_send_log
WHERE sent_at < datetime('now', '-7 days')"""
)
conn.commit()
finally:
if conn:
conn.close()
# ── A/B Testing ────────────────────────────────────────────────────────────────────
def create_ab_test(campaign_id, user_id, variants, test_size=20, metric="opens", declare_after_hours=24):
"""Create an A/B test for a campaign.
variants: list of {label, subject, body, split_percent}
Returns list of variant dicts.
"""
conn = None
try:
conn = get_db()
# Verify campaign ownership
campaign = conn.execute(
"SELECT id FROM campaigns WHERE id = ? AND user_id = ?",
(campaign_id, user_id),
).fetchone()
if not campaign:
return []
# Create variants
variant_ids = []
for v in variants:
cur = conn.execute(
"""INSERT INTO campaign_ab_variants
(campaign_id, variant_label, subject, body, split_percent)
VALUES (?, ?, ?, ?, ?)""",
(campaign_id, v["label"], v["subject"], v.get("body"), v.get("split_percent", 50)),
)
variant_ids.append(cur.lastrowid)
# Update campaign with A/B config
conn.execute(
"""UPDATE campaigns SET ab_test_metric = ?, ab_test_size = ?, ab_declare_after_hours = ?
WHERE id = ?""",
(metric, test_size, declare_after_hours, campaign_id),
)
conn.commit()
return [get_ab_variant(vi, user_id) for vi in variant_ids]
finally:
if conn:
conn.close()
def get_ab_variant(variant_id, user_id):
"""Get a single A/B variant. Returns dict or None."""
conn = None
try:
conn = get_db()
row = conn.execute(
"""SELECT v.* FROM campaign_ab_variants v
JOIN campaigns c ON v.campaign_id = c.id
WHERE v.id = ? AND c.user_id = ?""",
(variant_id, user_id),
).fetchone()
return dict(row) if row else None
finally:
if conn:
conn.close()
def get_ab_variants(campaign_id, user_id):
"""Get all A/B variants for a campaign. Returns list of dicts."""
conn = None
try:
conn = get_db()
rows = conn.execute(
"""SELECT v.* FROM campaign_ab_variants v
JOIN campaigns c ON v.campaign_id = c.id
WHERE v.campaign_id = ? AND c.user_id = ?
ORDER BY v.id""",
(campaign_id, user_id),
).fetchall()
return [dict(r) for r in rows]
finally:
if conn:
conn.close()
def get_ab_test_info(campaign_id, user_id):
"""Get A/B test info for a campaign. Returns dict with variants, metrics, winner."""
conn = None
try:
conn = get_db()
campaign = conn.execute(
"""SELECT ab_test_metric, ab_test_size, ab_declare_after_hours, ab_declared, ab_winner_variant_id
FROM campaigns WHERE id = ? AND user_id = ?""",
(campaign_id, user_id),
).fetchone()
if not campaign:
return None
campaign = dict(campaign)
variants = get_ab_variants(campaign_id, user_id)
return {
"metric": campaign["ab_test_metric"],
"test_size": campaign["ab_test_size"],
"declare_after_hours": campaign["ab_declare_after_hours"],
"declared": bool(campaign["ab_declared"]),
"winner_variant_id": campaign["ab_winner_variant_id"],
"variants": variants,
}
finally:
if conn:
conn.close()
def declare_ab_winner(campaign_id, user_id):
"""Declare A/B test winner based on configured metric. Returns winner variant dict."""
conn = None
try:
conn = get_db()
metric = conn.execute(
"SELECT ab_test_metric, ab_declared FROM campaigns WHERE id = ?",
(campaign_id,),
).fetchone()
if not metric or metric["ab_declared"]:
return None
metric = metric["ab_test_metric"]
metric_col_map = {
"opens": "opened_count",
"clicks": "clicked_count",
"responses": "response_count",
}
metric_col = metric_col_map.get(metric, "opened_count")
winner = conn.execute(
f"""SELECT * FROM campaign_ab_variants
WHERE campaign_id = ?
ORDER BY {metric_col} DESC
LIMIT 1""",
(campaign_id,),
).fetchone()
if not winner:
return None
winner = dict(winner)
winner_id = winner["id"]
# Mark winner
conn.execute(
"""UPDATE campaign_ab_variants SET is_winner = 1
WHERE campaign_id = ? AND id = ?""",
(campaign_id, winner_id),
)
conn.execute(
"""UPDATE campaigns SET ab_declared = 1, ab_winner_variant_id = ?
WHERE id = ?""",
(winner_id, campaign_id),
)
conn.commit()
return winner
finally:
if conn:
conn.close()
def assign_recipients_to_variants(campaign_id, user_id):
"""Assign pending recipients to A/B test variants based on split percentages.
Returns assignment count.
"""
import random
conn = None
try:
conn = get_db()
# Get variants
variants = conn.execute(
"SELECT * FROM campaign_ab_variants WHERE campaign_id = ? ORDER BY id",
(campaign_id,),
).fetchall()
if not variants:
return 0
variants = [dict(v) for v in variants]
total_split = sum(v["split_percent"] for v in variants)
# Get pending recipients
pending = conn.execute(
"""SELECT id FROM campaign_recipients
WHERE campaign_id = ? AND variant_id IS NULL AND status = 'pending'""",
(campaign_id,),
).fetchall()
if not pending:
return 0
# Build cumulative split map
assignments = 0
cumulative = 0
variant_map = []
for v in variants:
cumulative += v["split_percent"]
variant_map.append((cumulative, v["id"]))
for recipient in pending:
# Assign to variant based on weighted random
roll = random.randint(1, total_split)
for cum, variant_id in variant_map:
if roll <= cum:
conn.execute(
"UPDATE campaign_recipients SET variant_id = ? WHERE id = ?",
(variant_id, recipient["id"]),
)
assignments += 1
break
conn.commit()
return assignments
finally:
if conn:
conn.close()
def send_ab_winner_to_remaining(campaign_id, user_id):
"""Send the winning variant to recipients not yet sent.
Used after A/B test declares winner to promote to remaining recipients.
"""
from app.services.campaign_emails import send_campaign_email
conn = None
try:
conn = get_db()
winner_id = conn.execute(
"SELECT ab_winner_variant_id FROM campaigns WHERE id = ?",
(campaign_id,),
).fetchone()["ab_winner_variant_id"]
if not winner_id:
return 0
winner = conn.execute(
"SELECT * FROM campaign_ab_variants WHERE id = ?",
(winner_id,),
).fetchone()
if not winner:
return 0
winner = dict(winner)
# Find unsent recipients
pending = conn.execute(
"""SELECT * FROM campaign_recipients
WHERE campaign_id = ? AND status = 'pending'""",
(campaign_id,),
).fetchall()
sent = 0
for recipient in pending:
recipient = dict(recipient)
recipient["ab_variant"] = winner
success, _ = send_campaign_email(campaign_id, recipient)
if success:
sent += 1
return sent
finally:
if conn:
conn.close()
def check_pending_ab_tests():
"""Background job: check if any A/B tests are ready to declare winner.
Returns list of declared campaigns.
"""
from datetime import datetime, timedelta
conn = None
try:
conn = get_db()
now = datetime.now(UTC)
campaigns = conn.execute(
"""SELECT id, user_id, ab_declare_after_hours
FROM campaigns
WHERE ab_declared = 0
AND ab_test_metric IS NOT NULL""",
).fetchall()
declared = []
for campaign in campaigns:
campaign = dict(campaign)
hours = campaign.get("ab_declare_after_hours") or 24
threshold = now - timedelta(hours=hours)
# Check if campaign was sent long enough ago
sent_at = conn.execute(
"""SELECT MIN(sent_at) as sent_at FROM campaign_recipients
WHERE campaign_id = ? AND status = 'sent'""",
(campaign["id"],),
).fetchone()["sent_at"]
if not sent_at:
continue
try:
sent_time = datetime.fromisoformat(str(sent_at))
except (ValueError, TypeError):
continue
if sent_time < threshold:
winner = declare_ab_winner(campaign["id"], campaign["user_id"])
if winner:
declared.append(
{
"campaign_id": campaign["id"],
"winner": winner,
}
)
return declared
finally:
if conn:
conn.close()