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()