From c09470fe496f91debff7c246dd7721baa7530346 Mon Sep 17 00:00:00 2001 From: Vladless Date: Thu, 7 May 2026 17:09:48 +0000 Subject: [PATCH] Backend: broadcast channel select (bot/site/both) / scheduled_broadcasts.channel migration v29 / API v1+v2 channel field. Frontend: global haptic on button click / navigator.vibrate fallback for non-tg. --- api/v1/routes/management.py | 7 ++++ api/v2/routes/management.py | 7 ++++ database/migrations/schema_upgrade.py | 11 ++++++ database/models/notifications.py | 1 + database/scheduled_broadcasts.py | 2 + handlers/admin/sender/keyboard.py | 21 +++++++++++ handlers/admin/sender/scheduled_service.py | 12 +++++- handlers/admin/sender/sender_handler.py | 38 +++++++++++++++++-- handlers/admin/sender/sender_service.py | 43 +++++++++++++--------- handlers/admin/sender/sender_states.py | 1 + 10 files changed, 122 insertions(+), 21 deletions(-) diff --git a/api/v1/routes/management.py b/api/v1/routes/management.py index f868d9d6..2d4c9b13 100644 --- a/api/v1/routes/management.py +++ b/api/v1/routes/management.py @@ -58,6 +58,7 @@ class DomainChange(BaseModel): class BroadcastLaunchPayload(BaseModel): send_to: Literal["all", "subscribed", "unsubscribed", "untrial", "trial", "hotleads", "cluster"] = "all" + channel: Literal["bot", "site", "both"] = "both" text: str photo: str | None = None cluster_name: str | None = None @@ -71,6 +72,7 @@ class ScheduledBroadcastCreatePayload(BroadcastLaunchPayload): class ScheduledBroadcastUpdatePayload(BaseModel): send_to: Literal["all", "subscribed", "unsubscribed", "untrial", "trial", "hotleads", "cluster"] | None = None + channel: Literal["bot", "site", "both"] | None = None text: str | None = None photo: str | None = None cluster_name: str | None = None @@ -103,6 +105,7 @@ def _resolve_update_payload( fields = payload.model_fields_set text_changed = "text" in fields send_to = payload.send_to if "send_to" in fields else current.send_to + channel = payload.channel if "channel" in fields else current.channel text = payload.text if "text" in fields else current.text photo = payload.photo if "photo" in fields else current.photo cluster_name = payload.cluster_name if "cluster_name" in fields else current.cluster_name @@ -117,6 +120,7 @@ def _resolve_update_payload( cluster_name=cluster_name, workers=workers, messages_per_second=messages_per_second, + channel=channel, ) if not text_changed: prepared["text"] = current.text @@ -252,6 +256,7 @@ async def launch_broadcast( cluster_name=payload.cluster_name, workers=payload.workers, messages_per_second=payload.messages_per_second, + channel=payload.channel, ) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc @@ -272,6 +277,7 @@ async def create_broadcast_schedule( cluster_name=payload.cluster_name, workers=payload.workers, messages_per_second=payload.messages_per_second, + channel=payload.channel, ) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc @@ -279,6 +285,7 @@ async def create_broadcast_schedule( session, created_by_tg_id=getattr(admin, "tg_id", None), send_to=prepared["send_to"], + channel=prepared["channel"], cluster_name=prepared["cluster_name"], text=prepared["text"], photo=prepared["photo"], diff --git a/api/v2/routes/management.py b/api/v2/routes/management.py index 5cb1a17a..66ff0fce 100644 --- a/api/v2/routes/management.py +++ b/api/v2/routes/management.py @@ -73,6 +73,7 @@ class DomainChange(BaseModel): class BroadcastLaunchPayload(BaseModel): send_to: Literal["all", "subscribed", "unsubscribed", "untrial", "trial", "hotleads", "cluster"] = "all" + channel: Literal["bot", "site", "both"] = "both" text: str photo: str | None = None cluster_name: str | None = None @@ -86,6 +87,7 @@ class ScheduledBroadcastCreatePayload(BroadcastLaunchPayload): class ScheduledBroadcastUpdatePayload(BaseModel): send_to: Literal["all", "subscribed", "unsubscribed", "untrial", "trial", "hotleads", "cluster"] | None = None + channel: Literal["bot", "site", "both"] | None = None text: str | None = None photo: str | None = None cluster_name: str | None = None @@ -119,6 +121,7 @@ def _resolve_update_payload( fields = payload.model_fields_set text_changed = "text" in fields send_to = payload.send_to if "send_to" in fields else current.send_to + channel = payload.channel if "channel" in fields else current.channel text = payload.text if "text" in fields else current.text photo = payload.photo if "photo" in fields else current.photo cluster_name = payload.cluster_name if "cluster_name" in fields else current.cluster_name @@ -133,6 +136,7 @@ def _resolve_update_payload( cluster_name=cluster_name, workers=workers, messages_per_second=messages_per_second, + channel=channel, ) if not text_changed: prepared["text"] = current.text @@ -389,6 +393,7 @@ async def launch_broadcast( cluster_name=payload.cluster_name, workers=payload.workers, messages_per_second=payload.messages_per_second, + channel=payload.channel, ) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc @@ -409,6 +414,7 @@ async def create_broadcast_schedule( cluster_name=payload.cluster_name, workers=payload.workers, messages_per_second=payload.messages_per_second, + channel=payload.channel, ) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc @@ -416,6 +422,7 @@ async def create_broadcast_schedule( session, created_by_tg_id=getattr(identity, "tg_id", None), send_to=prepared["send_to"], + channel=prepared["channel"], cluster_name=prepared["cluster_name"], text=prepared["text"], photo=prepared["photo"], diff --git a/database/migrations/schema_upgrade.py b/database/migrations/schema_upgrade.py index 3d754b40..5ad5738e 100644 --- a/database/migrations/schema_upgrade.py +++ b/database/migrations/schema_upgrade.py @@ -1337,6 +1337,16 @@ async def _migration_v24_add_identity_sessions(conn: AsyncConnection) -> None: ) +async def _migration_v29_add_scheduled_broadcasts_channel(conn: AsyncConnection) -> None: + logger.info("[schema_upgrade] v29: scheduled_broadcasts.channel (bot/site/both)") + if not await _table_exists(conn, "scheduled_broadcasts"): + return + if not await _column_exists(conn, "scheduled_broadcasts", "channel"): + await conn.execute( + text("ALTER TABLE scheduled_broadcasts ADD COLUMN channel VARCHAR(8) NOT NULL DEFAULT 'both'") + ) + + _MIGRATIONS = [ (1, "Добавление users.id", _migration_v1_add_users_id), (2, "Добавление user_id колонок", _migration_v2_add_user_id_columns), @@ -1366,6 +1376,7 @@ _MIGRATIONS = [ (26, "индексы keys(expiry_time/server_id/tariff_id)", _migration_v26_add_keys_indexes), (27, "admins.permissions (JSONB per-admin permissions)", _migration_v27_add_admins_permissions), (28, "таблица identity_notif_prefs (toggle каналов)", _migration_v28_add_identity_notif_prefs), + (29, "scheduled_broadcasts.channel (bot/site/both)", _migration_v29_add_scheduled_broadcasts_channel), ] diff --git a/database/models/notifications.py b/database/models/notifications.py index db73363d..be93855b 100644 --- a/database/models/notifications.py +++ b/database/models/notifications.py @@ -39,6 +39,7 @@ class ScheduledBroadcast(DictLikeMixin, Base): created_by_tg_id = Column(BigInteger, ForeignKey("users.tg_id", ondelete="SET NULL"), nullable=True, index=True) status = Column(String(32), nullable=False, server_default=sql_text("'scheduled'"), index=True) send_to = Column(String(32), nullable=False, index=True) + channel = Column(String(8), nullable=False, server_default=sql_text("'both'")) cluster_name = Column(String, nullable=True) text = Column(Text, nullable=False) photo = Column(String, nullable=True) diff --git a/database/scheduled_broadcasts.py b/database/scheduled_broadcasts.py index 4b57b7b3..32dbc729 100644 --- a/database/scheduled_broadcasts.py +++ b/database/scheduled_broadcasts.py @@ -34,6 +34,7 @@ async def create_scheduled_broadcast( workers: int, messages_per_second: int, status: str = SCHEDULED_BROADCAST_STATUS_SCHEDULED, + channel: str = "both", ) -> ScheduledBroadcast: created_by_uid = None mirror_tg = created_by_tg_id @@ -46,6 +47,7 @@ async def create_scheduled_broadcast( created_by_user_id=created_by_uid, created_by_tg_id=mirror_tg, send_to=send_to, + channel=channel, cluster_name=cluster_name, text=text, photo=photo, diff --git a/handlers/admin/sender/keyboard.py b/handlers/admin/sender/keyboard.py index 485ee98e..7ea4255d 100644 --- a/handlers/admin/sender/keyboard.py +++ b/handlers/admin/sender/keyboard.py @@ -10,6 +10,10 @@ class AdminSenderCallback(CallbackData, prefix="admin_sender"): data: str | None = None +class AdminSenderChannelCallback(CallbackData, prefix="admin_sender_ch"): + channel: str + + class ScheduledBroadcastCallback(CallbackData, prefix="sb"): action: str broadcast_id: str = "0" @@ -87,6 +91,23 @@ def build_clusters_kb(clusters: list) -> InlineKeyboardMarkup: return builder.as_markup() +def build_channel_kb() -> InlineKeyboardMarkup: + builder = InlineKeyboardBuilder() + builder.row( + InlineKeyboardButton(text="📢 Везде (бот + сайт)", callback_data=AdminSenderChannelCallback(channel="both").pack()), + ) + builder.row( + InlineKeyboardButton(text="📲 Только бот", callback_data=AdminSenderChannelCallback(channel="bot").pack()), + InlineKeyboardButton(text="🌐 Только сайт", callback_data=AdminSenderChannelCallback(channel="site").pack()), + ) + builder.row(build_admin_back_btn()) + return builder.as_markup() + + +def channel_label(channel: str) -> str: + return {"both": "📢 везде", "bot": "📲 бот", "site": "🌐 сайт"}.get(channel, channel) + + def build_broadcast_preview_kb() -> InlineKeyboardMarkup: return InlineKeyboardMarkup( inline_keyboard=[ diff --git a/handlers/admin/sender/scheduled_service.py b/handlers/admin/sender/scheduled_service.py index 5cf9c818..7b51ffd5 100644 --- a/handlers/admin/sender/scheduled_service.py +++ b/handlers/admin/sender/scheduled_service.py @@ -72,10 +72,13 @@ def prepare_broadcast_payload( cluster_name: str | None = None, workers: int | None = None, messages_per_second: int | None = None, + channel: str = "both", ) -> dict: text_raw = (text or "").strip() if not text_raw: raise ValueError("Broadcast text is required") + if channel not in ("bot", "site", "both"): + raise ValueError("channel must be one of: bot, site, both") normalized_cluster_name = (cluster_name or "").strip() or None if send_to == "cluster" and not normalized_cluster_name: raise ValueError("Cluster name is required for cluster broadcast") @@ -86,6 +89,7 @@ def prepare_broadcast_payload( keyboard_json = keyboard.model_dump() if keyboard else None return { "send_to": send_to, + "channel": channel, "text": clean_text, "photo": photo, "cluster_name": normalized_cluster_name, @@ -101,6 +105,7 @@ def scheduled_broadcast_to_dict(broadcast: ScheduledBroadcast) -> dict: "created_by_tg_id": broadcast.created_by_tg_id, "status": broadcast.status, "send_to": broadcast.send_to, + "channel": broadcast.channel, "cluster_name": broadcast.cluster_name, "text": broadcast.text, "photo": broadcast.photo, @@ -154,7 +159,11 @@ async def execute_broadcast_payload(payload: dict, bot: Bot | None = None) -> di session=None, messages_per_second=clamp_broadcast_rate(payload.get("messages_per_second")), ) - stats = await broadcast_service.broadcast(messages, workers=clamp_broadcast_workers(payload.get("workers"))) + stats = await broadcast_service.broadcast( + messages, + workers=clamp_broadcast_workers(payload.get("workers")), + channel=payload.get("channel", "both"), + ) blocked_ids = stats.get("blocked_user_ids") or [] if blocked_ids: async with async_session_maker() as session: @@ -177,6 +186,7 @@ async def execute_broadcast_payload(payload: dict, bot: Bot | None = None) -> di async def execute_scheduled_broadcast(broadcast: ScheduledBroadcast, bot: Bot | None = None) -> dict: payload = { "send_to": broadcast.send_to, + "channel": broadcast.channel, "text": broadcast.text, "photo": broadcast.photo, "cluster_name": broadcast.cluster_name, diff --git a/handlers/admin/sender/sender_handler.py b/handlers/admin/sender/sender_handler.py index e63f03b9..f387ee17 100644 --- a/handlers/admin/sender/sender_handler.py +++ b/handlers/admin/sender/sender_handler.py @@ -29,12 +29,15 @@ from logger import logger from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb from .keyboard import ( AdminSenderCallback, + AdminSenderChannelCallback, ScheduledBroadcastCallback, build_broadcast_preview_kb, + build_channel_kb, build_clusters_kb, build_scheduled_broadcast_detail_kb, build_scheduled_broadcasts_list_kb, build_sender_kb, + channel_label, ) from .scheduled_service import ( execute_scheduled_broadcast, @@ -111,6 +114,7 @@ def _scheduled_broadcast_text(item) -> str: f"📌 Статус: {_scheduled_status_label(item.status)}", f"🕒 Время: {format_moscow_datetime(item.scheduled_for) or '-'}", f"👥 Аудитория: {target}", + f"📢 Канал: {channel_label(item.channel or 'both')}", f"🖼 Фото: {'Да' if item.photo else 'Нет'}", f"⌨️ Кнопки: {'Да' if item.keyboard_json else 'Нет'}", "", @@ -234,6 +238,7 @@ async def handle_broadcast_type( cluster_name=callback_data.data, workers=item.workers, messages_per_second=item.messages_per_second, + channel=item.channel, ) except ValueError as exc: await callback_query.message.edit_text(str(exc), reply_markup=build_admin_back_kb("sender")) @@ -257,12 +262,32 @@ async def handle_broadcast_type( reply_markup=build_scheduled_broadcast_detail_kb(updated, page=page), ) return + await callback_query.message.edit_text( + text="📢 Куда отправить рассылку?", + reply_markup=build_channel_kb(), + ) + await state.update_data(type=callback_data.type, cluster_name=callback_data.data) + await state.set_state(AdminSender.waiting_for_channel) + + +@router.callback_query( + AdminSenderChannelCallback.filter(), + AdminSender.waiting_for_channel, + IsAdminFilter(), +) +async def handle_channel_select( + callback_query: CallbackQuery, + callback_data: AdminSenderChannelCallback, + state: FSMContext, +): + if callback_data.channel not in ("bot", "site", "both"): + return + await state.update_data(channel=callback_data.channel) + await state.set_state(AdminSender.waiting_for_message) await callback_query.message.edit_text( text=_compose_message_text(), reply_markup=build_admin_back_kb("sender"), ) - await state.update_data(type=callback_data.type, cluster_name=callback_data.data) - await state.set_state(AdminSender.waiting_for_message) @router.message(AdminSender.waiting_for_message, IsAdminFilter()) @@ -284,6 +309,7 @@ async def handle_message_input(message: Message, state: FSMContext, session: Asy data = await state.get_data() send_to = data.get("type", "all") cluster_name = data.get("cluster_name") + channel = data.get("channel", "both") _, user_count = await get_recipients(session, send_to, cluster_name) if keyboard: @@ -310,7 +336,9 @@ async def handle_message_input(message: Message, state: FSMContext, session: Asy await message.answer(text=clean_text, parse_mode="HTML", reply_markup=keyboard) await message.answer( - f"👀 Это предпросмотр рассылки.\n👥 Количество получателей: {user_count}\n\nОтправить?", + f"👀 Это предпросмотр рассылки.\n" + f"👥 Количество получателей: {user_count}\n" + f"📢 Канал: {channel_label(channel)}\n\nОтправить?", reply_markup=build_broadcast_preview_kb(), ) @@ -323,6 +351,7 @@ async def handle_broadcast_confirm(callback_query: CallbackQuery, state: FSMCont keyboard_data = data.get("keyboard") send_to = data.get("type", "all") cluster_name = data.get("cluster_name") + channel = data.get("channel", "both") keyboard = None if keyboard_data: @@ -389,6 +418,7 @@ async def handle_broadcast_confirm(callback_query: CallbackQuery, state: FSMCont photo, state_keyboard_data, progress_cb, + channel, ) if stats.get("blocked_user_ids"): try: @@ -420,6 +450,7 @@ async def handle_broadcast_confirm(callback_query: CallbackQuery, state: FSMCont on_progress=on_progress, progress_interval=2.0, progress_every=200, + channel=channel, ) await callback_query.message.answer( @@ -461,6 +492,7 @@ async def handle_schedule_datetime_input(message: Message, state: FSMContext, se session, created_by_tg_id=message.from_user.id if message.from_user else None, send_to=data.get("type", "all"), + channel=data.get("channel", "both"), cluster_name=data.get("cluster_name"), text=data.get("text", ""), photo=data.get("photo"), diff --git a/handlers/admin/sender/sender_service.py b/handlers/admin/sender/sender_service.py index 7bb992cf..f42c0c03 100644 --- a/handlers/admin/sender/sender_service.py +++ b/handlers/admin/sender/sender_service.py @@ -28,10 +28,12 @@ def run_broadcast_in_thread( photo: str | None, keyboard_data: dict | None, progress_cb: Callable[[int, int, int, int, int], None] | None = None, + channel: str = "both", ) -> dict: """ Синхронная обёртка: запускает рассылку в отдельном event loop в текущем потоке. progress_cb принимает (completed, total, sent, failed, pending_retries). + channel: 'bot' / 'site' / 'both'. """ loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) @@ -52,6 +54,7 @@ def run_broadcast_in_thread( workers=5, on_progress=on_progress, progress_interval=2.0, + channel=channel, ) ) finally: @@ -282,7 +285,10 @@ class BroadcastService: on_progress: Callable[[int, int, int, int, int], Awaitable[None]] | None = None, progress_interval: float = 2.0, progress_every: int = 50, + channel: str = "both", ) -> dict: + send_to_bot = channel in ("bot", "both") + send_to_site = channel in ("site", "both") self.is_running = True self.start_time = time.time() self.results = [] @@ -290,18 +296,19 @@ class BroadcastService: self.blocked_users = set() self.pending_retries = 0 - for msg_data in messages: - msg = BroadcastMessage( - tg_id=msg_data["tg_id"], - text=msg_data["text"], - photo=msg_data.get("photo"), - keyboard=msg_data.get("keyboard"), - ) - await self.queue.put(msg) + if send_to_bot: + for msg_data in messages: + msg = BroadcastMessage( + tg_id=msg_data["tg_id"], + text=msg_data["text"], + photo=msg_data.get("photo"), + keyboard=msg_data.get("keyboard"), + ) + await self.queue.put(msg) - total = len(messages) + total = len(messages) if send_to_bot else 0 logger.info( - f"📤 Начата рассылка на {total} пользователей с {workers} воркерами " + f"📤 Начата рассылка channel={channel} на {len(messages)} получателей с {workers} воркерами " f"(rate={self.rate_limiter.max_rate}/сек, max_attempts={self.max_attempts})" ) @@ -311,13 +318,14 @@ class BroadcastService: self._progress_loop(total, on_progress, progress_interval, max(1, progress_every)), ) - worker_tasks = [asyncio.create_task(self._worker()) for _ in range(workers)] + worker_tasks = [asyncio.create_task(self._worker()) for _ in range(workers)] if send_to_bot else [] - while True: - await self.queue.join() - if self.pending_retries <= 0: - break - await asyncio.sleep(0.5) + if send_to_bot: + while True: + await self.queue.join() + if self.pending_retries <= 0: + break + await asyncio.sleep(0.5) self.is_running = False @@ -367,7 +375,8 @@ class BroadcastService: f"скорость: {avg_speed:.1f} сообщений/сек, время: {total_duration:.1f} сек" ) - await self._create_web_notifications(messages) + if send_to_site: + await self._create_web_notifications(messages) return stats diff --git a/handlers/admin/sender/sender_states.py b/handlers/admin/sender/sender_states.py index 42c6fedc..8d1677c8 100644 --- a/handlers/admin/sender/sender_states.py +++ b/handlers/admin/sender/sender_states.py @@ -2,6 +2,7 @@ from aiogram.fsm.state import State, StatesGroup class AdminSender(StatesGroup): + waiting_for_channel = State() waiting_for_message = State() preview = State() waiting_for_schedule_datetime = State()