#!/usr/bin/env python3 """ SuperSam — Paid Storage Campaign (State Machine) Уведомление о включении платного хранения. Логика: Когда delivery_status = 'paid_storage' (логист/админ включил): → Если рабочие часы (8-21, Пн-Пт): отправить SMS один раз → Если ночью: ждать до утра → Проверить доставку (103) → Пометить notification_status = 'paid_storage_sent' SMS текст: "Ваш заказ переведён на платное хранение. Стоимость: 300 ₽/день. Заберите заказ или согласуйте доставку: {link}" Защита от повторной отправки: 1. После отправки → notification_status = 'paid_storage_sending' 2. EXISTS проверка в sms_campaign_log за 24h """ import os import sys import logging import time import fcntl 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-paid-storage.log" CAMPAIGN_TYPE = "paid_storage" 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_paid_storage") # ─── БД ────────────────────────────────────────────────────────────────────── 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: return { "enabled": True, "test_mode": True, "test_phone": "79788382260", "max_check_duration_minutes": 90, "max_attempts": 2, "telegram_chat_id": TELEGRAM_CHAT_ID, "sms_api_id": SMS_API_ID, "work_hours_start": 8, "work_hours_end": 21, "work_days": "1,2,3,4,5", } return dict(row) # ─── SMS API ───────────────────────────────────────────────────────────────── def normalize_phone(phone): clean = "".join(c for c in str(phone) if c.isdigit()) if len(clean) < 10: return None 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): try: clean_phone = normalize_phone(phone) if not clean_phone: log.error(f"Invalid phone: {phone}") return None, "Invalid phone", "error" 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: return lines[1].strip(), text, "100" 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): 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): 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}") # ─── Проверка рабочего времени ──────────────────────────────────────────────── def is_within_work_hours(settings): now_msk = datetime.now(timezone(timedelta(hours=3))) 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 # ─── State Machine ─────────────────────────────────────────────────────────── def get_groups_to_send(conn): """Группы с paid_storage, которым ещё не отправили уведомление.""" 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, og.paid_storage_at FROM order_groups og WHERE og.delivery_status = 'paid_storage' AND og.paid_storage_at IS NOT NULL AND COALESCE(og.notification_status, '') NOT IN ('paid_storage_sending', 'paid_storage_sent') AND og.delivery_link IS NOT NULL AND og.delivery_link != '' AND NOT EXISTS ( SELECT 1 FROM sms_campaign_log scl WHERE scl.order_group_id = og.id AND scl.campaign_type = 'paid_storage' AND scl.status IN ('sent', 'checking') AND scl.created_at > NOW() - INTERVAL '24 hours' ) ORDER BY og.paid_storage_at ASC """) return [dict(r) for r in cur.fetchall()] def get_sms_to_check(conn, max_duration_min): 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 = 'paid_storage' 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): 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 = 'paid_storage' 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 step_send_new(conn, settings, test_send=False): api_id = settings.get("sms_api_id", SMS_API_ID) tg_chat = settings.get("telegram_chat_id", TELEGRAM_CHAT_ID) if test_send: # Test mode: send ONE SMS to test phone, pick first available group all_groups = get_groups_to_send(conn) if not all_groups: log.info("Test send: no groups available in queue") return 0 groups = [all_groups[0]] # Only first group log.info(f"TEST SEND: sending 1 SMS to test phone (skipping {len(all_groups)-1} others)") else: groups = get_groups_to_send(conn) log.info(f"Step 1: {len(groups)} groups to send paid_storage 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", "") with conn.cursor() as cur: cur.execute(""" SELECT COUNT(*) FROM sms_campaign_log WHERE order_group_id = %s AND campaign_type = 'paid_storage' AND created_at > NOW() - INTERVAL '24 hours' AND status IN ('sent', 'checking', 'delivered') """, (group_id,)) if cur.fetchone()[0] > 0: log.info(f"Group {group_id}: already has recent paid_storage SMS, skipping") continue sms_text_template = settings.get("sms_text_template", "Ваш заказ переведён на платное хранение. Стоимость: 300 ₽/день. Заберите заказ или согласуйте доставку: {link}") sms_text = sms_text_template.replace("{link}", 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 paid_storage 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, sent_to=send_phone, status="sent", sms_code=code, attempts=1, was_test_mode=bool(settings.get("test_mode", True)), ) update_order_group(conn, group_id, { "notification_status": "paid_storage_sending", "sms_sent_at": "NOW()", }) log.info(f"Group {group_id}: paid_storage SMS sent, sms_id={sms_id}, notification_status→paid_storage_sending") sent_count += 1 time.sleep(settings.get("send_interval_seconds", 15)) # Configurable rate limit else: error = raw[:500] if raw else "Unknown error" insert_sms_log(conn, campaign_type=CAMPAIGN_TYPE, order_group_id=group_id, customer_phone=phone, sms_text=sms_text, sent_to=send_phone if "send_phone" in dir() else phone, status="send_failed", sms_code=code, attempts=1, error_message=error, was_test_mode=bool(settings.get("test_mode", True)), ) log.error(f"Group {group_id}: paid_storage SMS failed (code={code}): {error[:200]}") update_order_group(conn, group_id, {"last_sms_error": error[:200]}) send_telegram(f"❌ SMS платное хранение не отправлена: {name} ({phone})\nКод: {code}", tg_chat) return sent_count def step_check_status(conn, settings): 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) sms_list = get_sms_to_check(conn, max_duration) log.info(f"Step 2: {len(sms_list)} paid_storage SMS to check") 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", "") code, raw, api_code = check_sms_status(sms_id, api_id) update_sms_log(conn, log_id, needs_check=False, checked_at="NOW()") if code == DELIVERED_CODE: log.info(f"Group {group_id}: paid_storage SMS delivered (103)!") update_sms_log(conn, log_id, status="delivered", sms_code=code) update_order_group(conn, group_id, {"notification_status": "paid_storage_sent"}) send_telegram(f"✅ SMS платное хранение доставлена: {name} ({phone})", tg_chat) delivered += 1 elif code in IN_TRANSIT_CODES: log.info(f"Group {group_id}: in transit (code={code})") update_sms_log(conn, log_id, status="checking", sms_code=code) elif code in DELIVERY_ERROR_CODES: log.error(f"Group {group_id}: delivery error (code={code})") update_sms_log(conn, log_id, status="error", sms_code=code, error_message=f"Delivery error: {code}") attempts = item.get("attempts", 1) max_attempts = settings.get("max_attempts", 2) if attempts < max_attempts: update_sms_log(conn, log_id, status="expired") # Сброс для retry — вернёмся к исходному статусу update_order_group(conn, group_id, {"notification_status": "not_started"}) log.info(f"Group {group_id}: will retry paid_storage SMS (attempt {attempts+1}/{max_attempts})") else: update_order_group(conn, group_id, { "notification_status": "manual_required", "last_sms_error": f"Paid storage SMS failed after {max_attempts} attempts (code={code})", }) send_telegram(f"⚠️ SMS платное хранение не доставлена после {max_attempts} попыток: {name} ({phone})", tg_chat) elif code in LIMIT_ERROR_CODES: log.error(f"Group {group_id}: limit exceeded (code={code})") update_sms_log(conn, log_id, status="limit_exceeded", sms_code=code, error_message=f"Limit: {code}") update_order_group(conn, group_id, { "notification_status": "manual_required", "last_sms_error": f"Limit exceeded (code={code})", }) send_telegram(f"🚫 SMS платное хранение заблокирована (лимит): {name} ({phone})", tg_chat) else: log.warning(f"Group {group_id}: unknown code: {code}") update_sms_log(conn, log_id, status="checking", sms_code=code, error_message=f"Unknown: {code}") return delivered def step_handle_expired(conn, settings): max_duration = settings.get("max_check_duration_minutes", 90) tg_chat = settings.get("telegram_chat_id", TELEGRAM_CHAT_ID) expired = get_sms_expired(conn, max_duration) log.info(f"Step 3: {len(expired)} paid_storage SMS expired") 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", "") log.warning(f"Group {group_id}: paid_storage SMS expired") update_sms_log(conn, log_id, status="expired", error_message=f"Not delivered in {max_duration} min") attempts = item.get("attempts", 1) max_attempts = settings.get("max_attempts", 2) if attempts < max_attempts: # Retry: сброс на not_started для повторной отправки update_order_group(conn, group_id, {"notification_status": "not_started"}) log.info(f"Group {group_id}: will retry paid_storage SMS (attempt {attempts+1}/{max_attempts})") else: # Все попытки исчерпаны — ручное управление, НЕ paid_storage_sent update_order_group(conn, group_id, { "notification_status": "manual_required", "last_sms_error": f"Paid storage SMS not delivered after {max_attempts} attempts", }) send_telegram(f"⚠️ SMS платное хранение не доставлена после {max_attempts} попыток: {name} ({phone})", tg_chat) # ─── Main ──────────────────────────────────────────────────────────────────── def main(): # File lock — prevent parallel execution lock_file = open("/tmp/" + __file__.split("/")[-1].replace(".py", ".lock"), "w") try: fcntl.flock(lock_file, fcntl.LOCK_EX | fcntl.LOCK_NB) except (IOError, OSError): log.info("Another instance is running, exiting") lock_file.close() return log.info("=" * 60) log.info("Paid Storage Campaign — START") conn = get_db_conn() try: settings = load_settings(conn) # Update last_run_at timestamp with conn.cursor() as cur: cur.execute("UPDATE sms_campaign_settings SET last_run_at = NOW() WHERE campaign_type = %s", (CAMPAIGN_TYPE,)) conn.commit() # Check test_send_requested — one-time test send test_send_requested = bool(settings.get("test_send_requested", False)) if test_send_requested: # Reset flag immediately with conn.cursor() as cur: cur.execute("UPDATE sms_campaign_settings SET test_send_requested = false WHERE campaign_type = %s", (CAMPAIGN_TYPE,)) conn.commit() log.info("Test send requested — will send ONE SMS to test phone only") # Check run_requested — restart scenario run_requested = bool(settings.get('run_requested', False)) if run_requested: with conn.cursor() as cur: cur.execute("UPDATE sms_campaign_settings SET run_requested = false WHERE campaign_type = 'paid_storage'") conn.commit() with conn.cursor() as cur: cur.execute("""UPDATE order_groups SET notification_status = 'not_started', sms_sent_at = NULL, next_notification_check_at = NULL, sms_attempts = 0, last_sms_error = NULL WHERE delivery_status = 'paid_storage' AND notification_status IN ('paid_storage_sent','paid_storage_sending','manual_required','send_failed')""") reset_count = cur.rowcount conn.commit() log.info(f'RESTART paid_storage: reset {reset_count} groups') log.info(f"Settings: work={settings.get('work_hours_start')}-{settings.get('work_hours_end')}, " f"days={settings.get('work_days')}, test={settings.get('test_mode')}") if not settings.get("enabled", True): log.info("Campaign disabled, exiting") return # Отправка только в рабочие часы work_hours = is_within_work_hours(settings) sent = 0 if work_hours: sent = step_send_new(conn, settings, test_send=test_send_requested) else: log.info("Outside work hours, 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"📦 Платное хранение\nОтправлено: {sent}\nДоставлено: {delivered}", settings.get("telegram_chat_id", TELEGRAM_CHAT_ID), was_test_mode=bool(settings.get("test_mode", True)), ) # Обновляем баланс 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 = 'paid_storage'", (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() fcntl.flock(lock_file, fcntl.LOCK_UN) lock_file.close() log.info("Paid Storage Campaign — END") log.info("=" * 60) if __name__ == "__main__": main()