supersam/scripts/sms_second_campaign.py

587 lines
27 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""
SuperSam — SMS Second Campaign (State Machine)
Заменяет n8n workflow "Первая проверка согласования + отправка 2й смс"
Архитектура: state machine через cron (каждые 5 мин)
Каждый запуск:
1. Отправляет вторую SMS группам, где первая SMS доставлена,
но клиент не согласовал доставку, и пришло время next_notification_check_at
2. Проверяет статус ранее отправленных вторых SMS
3. Обновляет статусы в order_groups + sms_campaign_log
Защита от повторной отправки (ДВОЙНАЯ):
1. После отправки SMS → notification_status = 'second_sms_sending' (не 'first_sms_sent')
→ get_groups_to_send НЕ находит эту группу
2. Дополнительно: EXISTS проверка в sms_campaign_log (status sent/checking за последние 24h)
3. Код 231/132 = лимит одинаковых → не повторять
needs_check=true — admin нажал "Проверить снова" в UI
→ скрипт проверяет даже если запись старше max_check_duration
Коды sms.ru:
100 = в очереди, 101 = оператору, 102 = в пути → ждём
103 = доставлено → цель
104-108, 130-132, 230-232 = ошибки доставки / лимиты
"""
import os
import sys
import json
import logging
from datetime import datetime, timezone, timedelta
import requests
import psycopg2
from psycopg2.extras import RealDictCursor
# ─── Конфигурация ────────────────────────────────────────────────────────────
DB_HOST = os.environ.get("DB_HOST", "10.0.4.12")
DB_PORT = os.environ.get("DB_PORT", "5432")
DB_NAME = os.environ.get("DB_NAME", "postgres")
DB_USER = os.environ.get("DB_USER", "supabase_admin")
DB_PASS = os.environ.get("DB_PASS", "4fe80bb21c7c3d17a8d8b226adf7a479")
TELEGRAM_BOT_TOKEN = os.environ.get("TELEGRAM_BOT_TOKEN", "")
TELEGRAM_CHAT_ID = os.environ.get("TELEGRAM_CHAT_ID", "25164483")
SMS_API_ID = os.environ.get("SMS_API_ID", "C92063B3-95ED-8559-157B-1946EB5A2486")
SMS_SEND_URL = "https://sms.ru/sms/send"
SMS_STATUS_URL = "https://sms.ru/sms/status"
LOG_FILE = "/var/log/supersam-sms-second.log"
CAMPAIGN_TYPE = "second_sms"
# Коды, которые означают "в процессе" (ждём дальше)
IN_TRANSIT_CODES = {"100", "101", "102"}
# Код доставки
DELIVERED_CODE = "103"
# Коды ошибок доставки (не временные, можно повторять отправку)
DELIVERY_ERROR_CODES = {"104", "105", "106", "107", "108", "130"}
# Коды превышения лимитов (нельзя повторять отправку — заблокируют)
LIMIT_ERROR_CODES = {"131", "132", "230", "231", "232"}
# ─── Логирование ─────────────────────────────────────────────────────────────
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
handlers=[
logging.FileHandler(LOG_FILE),
logging.StreamHandler(sys.stdout),
],
)
log = logging.getLogger("sms_second")
# ─── БД ──────────────────────────────────────────────────────────────────────
def get_db_conn():
return psycopg2.connect(
host=DB_HOST, port=DB_PORT, dbname=DB_NAME,
user=DB_USER, password=DB_PASS,
)
def load_settings(conn):
with conn.cursor(cursor_factory=RealDictCursor) as cur:
cur.execute("SELECT * FROM sms_campaign_settings WHERE campaign_type = %s", (CAMPAIGN_TYPE,))
row = cur.fetchone()
if not row:
# Fallback на настройки first_sms если second_sms нет
cur.execute("SELECT * FROM sms_campaign_settings WHERE campaign_type = 'first_sms'")
row = cur.fetchone()
if not row:
return {
"wait_between_checks_seconds": 25,
"max_check_duration_minutes": 90,
"max_attempts": 2,
"enabled": True,
"telegram_chat_id": "25164483",
"sms_api_id": SMS_API_ID,
"test_mode": True,
"test_phone": "79788382260",
"second_sms_delay_hours": 3,
"auto_manual_after_hours": 3,
}
return dict(row)
# ─── SMS API ─────────────────────────────────────────────────────────────────
def normalize_phone(phone):
"""Нормализует телефон: только цифры, начинается с 7"""
clean = "".join(c for c in str(phone) if c.isdigit())
if clean.startswith("8"):
clean = "7" + clean[1:]
elif not clean.startswith("7"):
clean = "7" + clean
return clean
def send_sms(phone, message, api_id):
"""Отправляет SMS, возвращает (sms_id, raw_response, code)"""
try:
clean_phone = normalize_phone(phone)
resp = requests.post(SMS_SEND_URL, params={
"api_id": api_id,
"to": clean_phone,
}, data={"msg": message}, timeout=30)
text = resp.text
log.info(f"SMS send to {clean_phone}: {text[:200]}")
lines = text.strip().split("\n")
status_code = lines[0].strip() if lines else ""
if status_code == "100" and len(lines) >= 2:
sms_id = lines[1].strip()
return sms_id, text, "100"
else:
return None, text, status_code
except Exception as e:
log.error(f"SMS send error: {e}")
return None, str(e), "error"
def check_sms_status(sms_id, api_id):
"""Проверяет статус, возвращает (sms_status_code, raw_response, api_code)"""
try:
resp = requests.post(SMS_STATUS_URL, params={
"api_id": api_id,
"sms_id": sms_id,
}, timeout=30)
text = resp.text
lines = text.strip().split("\n")
status_code = lines[0].strip() if lines else ""
sms_status_code = lines[1].strip() if len(lines) >= 2 else None
log.info(f"SMS status for {sms_id}: code={status_code}, sms_status={sms_status_code}")
return sms_status_code, text, status_code
except Exception as e:
log.error(f"SMS status check error: {e}")
return None, str(e), "error"
def fetch_balance(api_id):
"""Получает баланс sms.ru, возвращает (balance_float, raw)"""
try:
resp = requests.get("https://sms.ru/my/balance", params={"api_id": api_id}, timeout=15)
text = resp.text
lines = text.strip().split("\n")
if lines[0].strip() == "100" and len(lines) >= 2:
return float(lines[1].strip()), text
return None, text
except Exception as e:
log.error(f"Balance fetch error: {e}")
return None, str(e)
# ─── Telegram ────────────────────────────────────────────────────────────────
def send_telegram(message, chat_id):
if not TELEGRAM_BOT_TOKEN:
log.warning("TELEGRAM_BOT_TOKEN not set, skipping Telegram")
return
try:
resp = requests.post(
f"https://api.telegram.org/bot{TELEGRAM_BOT_TOKEN}/sendMessage",
json={"chat_id": chat_id, "text": message, "parse_mode": "HTML"},
timeout=10,
)
if resp.status_code != 200:
log.warning(f"Telegram error: {resp.text[:200]}")
except Exception as e:
log.warning(f"Telegram failed: {e}")
# ─── State Machine ───────────────────────────────────────────────────────────
def get_groups_to_send(conn):
"""Группы, где первая SMS доставлена, но клиент не согласовал,
и пришло время для второй SMS (next_notification_check_at <= NOW()).
second_sms_sent_at IS NULL — вторая SMS ещё не отправлена.
"""
with conn.cursor(cursor_factory=RealDictCursor) as cur:
cur.execute("""
SELECT og.id, og.group_key, og.customer_name, og.customer_phone,
og.customer_phone_normalized, og.delivery_link, og.notification_status
FROM order_groups og
WHERE og.delivery_status = 'pending_confirmation'
AND og.delivery_link IS NOT NULL
AND og.delivery_link != ''
AND COALESCE(og.notification_status, '') = 'first_sms_sent'
AND og.second_sms_sent_at IS NULL
AND (og.next_notification_check_at IS NULL OR og.next_notification_check_at <= NOW())
AND sms_allowed_for_city(og.town, 'second_sms')
-- Нет активной второй SMS в логе
AND NOT EXISTS (
SELECT 1 FROM sms_campaign_log scl
WHERE scl.order_group_id = og.id
AND scl.campaign_type = 'second_sms'
AND scl.status IN ('sent', 'checking')
AND scl.created_at > NOW() - INTERVAL '2 hours'
)
ORDER BY og.created_at ASC
""")
return [dict(r) for r in cur.fetchall()]
def get_sms_to_check(conn, max_duration_min):
"""Вторые SMS в логе со status='sent'/'checking', которые ещё не доставлены.
Включает:
- Записи младше max_duration_min (обычная автопроверка)
- Записи с needs_check=true (admin нажал 'Проверить снова') — независимо от возраста
"""
with conn.cursor(cursor_factory=RealDictCursor) as cur:
cur.execute("""
SELECT scl.id as log_id, scl.sms_id, scl.order_group_id, scl.customer_phone,
scl.attempts, scl.created_at, scl.sms_code, scl.needs_check,
og.customer_name, og.group_key
FROM sms_campaign_log scl
JOIN order_groups og ON og.id = scl.order_group_id
WHERE scl.campaign_type = 'second_sms'
AND scl.status IN ('sent', 'checking')
AND scl.sms_id IS NOT NULL
AND (
scl.created_at > NOW() - INTERVAL '%s minutes'
OR scl.needs_check = true
)
ORDER BY scl.needs_check DESC, scl.created_at ASC
""" % max_duration_min)
return [dict(r) for r in cur.fetchall()]
def get_sms_expired(conn, max_duration_min):
"""Вторые SMS, у которых истёк срок проверки (старше max_check_duration, не доставлены)"""
with conn.cursor(cursor_factory=RealDictCursor) as cur:
cur.execute("""
SELECT scl.id as log_id, scl.sms_id, scl.order_group_id, scl.customer_phone,
scl.attempts, scl.created_at,
og.customer_name, og.group_key, og.delivery_link
FROM sms_campaign_log scl
JOIN order_groups og ON og.id = scl.order_group_id
WHERE scl.campaign_type = 'second_sms'
AND scl.status IN ('sent', 'checking')
AND scl.sms_id IS NOT NULL
AND scl.created_at < NOW() - INTERVAL '%s minutes'
ORDER BY scl.created_at ASC
""" % max_duration_min)
return [dict(r) for r in cur.fetchall()]
def insert_sms_log(conn, **kwargs):
with conn.cursor() as cur:
cols = ", ".join(kwargs.keys())
placeholders = ", ".join(["%s"] * len(kwargs))
cur.execute(f"INSERT INTO sms_campaign_log ({cols}) VALUES ({placeholders}) RETURNING id", list(kwargs.values()))
log_id = cur.fetchone()[0]
conn.commit()
return log_id
def update_sms_log(conn, log_id, **kwargs):
with conn.cursor() as cur:
set_parts = []
values = []
for k, v in kwargs.items():
if v == "NOW()":
set_parts.append(f"{k} = NOW()")
else:
set_parts.append(f"{k} = %s")
values.append(v)
values.append(log_id)
cur.execute(f"UPDATE sms_campaign_log SET {', '.join(set_parts)}, updated_at = NOW() WHERE id = %s", values)
conn.commit()
def update_order_group(conn, group_id, fields):
with conn.cursor() as cur:
set_parts = []
values = []
for k, v in fields.items():
if v == "NOW()":
set_parts.append(f"{k} = NOW()")
else:
set_parts.append(f"{k} = %s")
values.append(v)
values.append(group_id)
cur.execute(f"UPDATE order_groups SET {', '.join(set_parts)} WHERE id = %s", values)
conn.commit()
# ─── Проверка рабочего времени ────────────────────────────────────────────────
def is_within_work_hours(settings):
"""Проверка: сейчас рабочие часы.
settings: work_hours_start, work_hours_end (часы 0-23), work_days ('1,2,3,4,5')
"""
now_msk = datetime.now(timezone(timedelta(hours=3)))
# work_days: '1,2,3,4,5' → Понедельник=1 ... Воскресенье=7
# Python weekday(): 0=Пн ... 6=Вс → конвертируем в 1-7
today_num = now_msk.weekday() + 1
allowed_days = set()
work_days_str = settings.get("work_days", "1,2,3,4,5")
for part in str(work_days_str).split(","):
part = part.strip()
if part.isdigit():
allowed_days.add(int(part))
if today_num not in allowed_days:
return False
hour = now_msk.hour
start_h = settings.get("work_hours_start", 8)
end_h = settings.get("work_hours_end", 21)
return start_h <= hour < end_h
# ─── Основная логика ─────────────────────────────────────────────────────────
def step_send_new(conn, settings):
"""Шаг 1: Отправка второй SMS группам, где первая доставлена, но нет согласования"""
api_id = settings.get("sms_api_id", SMS_API_ID)
tg_chat = settings.get("telegram_chat_id", TELEGRAM_CHAT_ID)
groups = get_groups_to_send(conn)
log.info(f"Step 1: {len(groups)} groups to send second SMS")
sent_count = 0
for group in groups:
group_id = str(group["id"])
name = group.get("customer_name") or group.get("group_key", "")
phone = group.get("customer_phone_normalized") or group.get("customer_phone", "")
delivery_link = group.get("delivery_link", "")
# Защита: была ли уже вторая SMS этой группе за 24h
with conn.cursor() as cur:
cur.execute("""
SELECT COUNT(*) as cnt FROM sms_campaign_log
WHERE order_group_id = %s AND campaign_type = 'second_sms'
AND created_at > NOW() - INTERVAL '24 hours'
AND status IN ('sent', 'checking', 'delivered')
""", (group_id,))
recent_count = cur.fetchone()[0]
if recent_count > 0:
log.info(f"Group {group_id}: already has recent second SMS in log, skipping")
continue
sms_text = f"Ваш заказ готов к доставке. Выберите дату доставки по ссылке: {delivery_link}"
# ТЕСТОВЫЙ РЕЖИМ
send_phone = phone
if settings.get("test_mode", True):
send_phone = settings.get("test_phone", "79788382260")
log.info(f"TEST MODE: sending to {send_phone} instead of {phone}")
log.info(f"Sending second SMS to {name} (orig={phone}, send={send_phone})")
sms_id, raw, code = send_sms(send_phone, sms_text, api_id)
if sms_id:
log_id = insert_sms_log(conn,
campaign_type=CAMPAIGN_TYPE,
order_group_id=group_id,
customer_phone=phone,
sms_id=sms_id,
sms_text=sms_text,
status="sent",
sms_code=code,
attempts=1,
)
# ДВОЙНАЯ ЗАЩИТА: notification_status → second_sms_sending
update_order_group(conn, group_id, {
"notification_status": "second_sms_sending",
"sms_sent_at": "NOW()",
})
log.info(f"Group {group_id}: second SMS sent, sms_id={sms_id}, log_id={log_id}, notification_status→second_sms_sending")
sent_count += 1
else:
error = raw[:500] if raw else "Unknown error"
log_id = insert_sms_log(conn,
campaign_type=CAMPAIGN_TYPE,
order_group_id=group_id,
customer_phone=phone,
sms_text=sms_text,
status="send_failed",
sms_code=code,
attempts=1,
error_message=error,
)
log.error(f"Group {group_id}: second SMS send failed (code={code}): {error[:200]}")
update_order_group(conn, group_id, {
"last_sms_error": error[:200],
})
send_telegram(f"❌ Вторая SMS не отправлена: {name} ({phone})\nКод: {code}\nОшибка: {error[:200]}", tg_chat)
return sent_count
def step_check_status(conn, settings):
"""Шаг 2: Проверка статуса ранее отправленных вторых SMS"""
api_id = settings.get("sms_api_id", SMS_API_ID)
tg_chat = settings.get("telegram_chat_id", TELEGRAM_CHAT_ID)
max_duration = settings.get("max_check_duration_minutes", 90)
max_attempts = settings.get("max_attempts", 2)
auto_manual_hours = settings.get("auto_manual_after_hours", 3)
sms_list = get_sms_to_check(conn, max_duration)
log.info(f"Step 2: {len(sms_list)} second SMS to check status")
delivered = 0
for item in sms_list:
log_id = item["log_id"]
sms_id = item["sms_id"]
group_id = str(item["order_group_id"])
name = item.get("customer_name") or item.get("group_key", "")
phone = item.get("customer_phone", "")
attempts = item.get("attempts", 1)
code, raw, api_code = check_sms_status(sms_id, api_id)
# Сбрасываем needs_check + ставим checked_at
update_sms_log(conn, log_id, needs_check=False, checked_at="NOW()")
if code == DELIVERED_CODE:
# Доставлено!
log.info(f"Group {group_id}: second SMS delivered (103)!")
update_sms_log(conn, log_id, status="delivered", sms_code=code)
next_check = (datetime.now(timezone.utc) + timedelta(hours=auto_manual_hours)).isoformat()
update_order_group(conn, group_id, {
"notification_status": "second_sms_sent",
"second_sms_sent_at": "NOW()",
"sms_attempts": attempts,
"last_sms_error": None,
"sms_sent_at": "NOW()",
"next_notification_check_at": next_check,
"status": "second_sms_sent",
})
send_telegram(f"✅ Вторая SMS доставлена: {name} ({phone})", tg_chat)
delivered += 1
elif code in IN_TRANSIT_CODES:
# В пути / в очереди — продолжаем ждать
log.info(f"Group {group_id}: second SMS in transit (code={code}), will check again next run")
update_sms_log(conn, log_id, status="checking", sms_code=code)
elif code in DELIVERY_ERROR_CODES:
# Ошибка доставки — можно повторить отправку
log.error(f"Group {group_id}: second SMS delivery error (code={code})")
update_sms_log(conn, log_id, status="error", sms_code=code,
error_message=f"Delivery error: {code}")
if attempts < max_attempts:
log.info(f"Group {group_id}: will retry second SMS send (attempt {attempts+1}/{max_attempts})")
update_sms_log(conn, log_id, status="expired")
# Сбрасываем на first_sms_sent для повторной отправки
update_order_group(conn, group_id, {
"notification_status": "first_sms_sent",
})
else:
update_order_group(conn, group_id, {
"notification_status": "manual_required",
"last_sms_error": f"Second SMS failed after {max_attempts} attempts (code={code})",
})
send_telegram(f"⚠️ Вторая SMS не доставлена после {max_attempts} попыток: {name} ({phone})\nКод: {code}", tg_chat)
elif code in LIMIT_ERROR_CODES:
# Превышен лимит — НЕ повторять
log.error(f"Group {group_id}: second SMS limit exceeded (code={code}) — NOT retrying")
update_sms_log(conn, log_id, status="limit_exceeded", sms_code=code,
error_message=f"Limit exceeded: code={code}")
update_order_group(conn, group_id, {
"notification_status": "manual_required",
"last_sms_error": f"Second SMS limit exceeded (code={code})",
})
send_telegram(f"🚫 Вторая SMS заблокирована (лимит {code}): {name} ({phone})", tg_chat)
else:
# Неизвестный код — логируем, продолжаем проверять
log.warning(f"Group {group_id}: unknown SMS code for second SMS: {code}")
update_sms_log(conn, log_id, status="checking", sms_code=code,
error_message=f"Unknown code: {code}")
return delivered
def step_handle_expired(conn, settings):
"""Шаг 3: Обработка вторых SMS с истёкшим сроком проверки"""
max_duration = settings.get("max_check_duration_minutes", 90)
max_attempts = settings.get("max_attempts", 2)
tg_chat = settings.get("telegram_chat_id", TELEGRAM_CHAT_ID)
expired = get_sms_expired(conn, max_duration)
log.info(f"Step 3: {len(expired)} second SMS expired (older than {max_duration} min)")
for item in expired:
log_id = item["log_id"]
group_id = str(item["order_group_id"])
name = item.get("customer_name") or item.get("group_key", "")
phone = item.get("customer_phone", "")
attempts = item.get("attempts", 1)
log.warning(f"Group {group_id}: second SMS expired after {max_duration} min, attempts={attempts}/{max_attempts}")
update_sms_log(conn, log_id, status="expired",
error_message=f"Not delivered in {max_duration} minutes")
if attempts < max_attempts:
# Сбрасываем на first_sms_sent для повторной отправки
update_order_group(conn, group_id, {
"notification_status": "first_sms_sent",
})
log.info(f"Group {group_id}: will retry second SMS send (attempt {attempts+1}), notification_status→first_sms_sent")
else:
# Все попытки исчерпаны → ручное управление
update_order_group(conn, group_id, {
"notification_status": "manual_required",
"last_sms_error": f"Second SMS not delivered after {max_attempts} attempts",
})
send_telegram(
f"🔧 Требуется ручное управление (2-я SMS): {name} ({phone})\n"
f"Вторая SMS не доставлена после {max_attempts} попыток",
tg_chat,
)
# ─── Main ────────────────────────────────────────────────────────────────────
def main():
log.info("=" * 60)
log.info("SMS Second Campaign — START")
conn = get_db_conn()
try:
settings = load_settings(conn)
log.info(f"Settings: wait={settings.get('wait_between_checks_seconds')}s, "
f"max_duration={settings.get('max_check_duration_minutes')}min, "
f"max_attempts={settings.get('max_attempts')}")
if not settings.get("enabled", True):
log.info("Campaign disabled, exiting")
return
# State machine — каждый шаг быстрый, без blocking
# Отправка только в рабочие часы (8-21, Пн-Пт по Москве)
work_hours = is_within_work_hours(settings)
sent = 0
if work_hours:
sent = step_send_new(conn, settings)
else:
log.info("Outside work hours (8-21 MSK, Mon-Fri), skipping new SMS sends")
# Проверка статусов работает всегда — даже ночью
delivered = step_check_status(conn, settings)
step_handle_expired(conn, settings)
log.info(f"Run summary: sent={sent}, delivered={delivered}")
if sent > 0 or delivered > 0:
send_telegram(
f"📊 <b>Вторая отправка</b>\nОтправлено: {sent}\nДоставлено: {delivered}",
settings.get("telegram_chat_id", TELEGRAM_CHAT_ID),
)
# Обновляем баланс sms.ru
if sent > 0:
api_id = settings.get("sms_api_id", SMS_API_ID)
balance, raw = fetch_balance(api_id)
if balance is not None:
with conn.cursor() as cur:
cur.execute(
"UPDATE sms_campaign_settings SET last_balance = %s WHERE campaign_type = 'second_sms'",
(balance,)
)
conn.commit()
log.info(f"Balance updated: {balance}")
except Exception as e:
log.error(f"Fatal error: {e}", exc_info=True)
finally:
conn.close()
log.info("SMS Second Campaign — END")
log.info("=" * 60)
if __name__ == "__main__":
main()