fix
This commit is contained in:
+46
-10
@@ -2091,8 +2091,29 @@ class BotInterface:
|
||||
reply_markup=builder.as_markup()
|
||||
)
|
||||
|
||||
async def _drop_pending_series(self, params: TaskParams, url_tasks: Optional[set]) -> int:
|
||||
"""Сбрасывает очередь кликов, накопленную за прошлую трансляцию.
|
||||
|
||||
Отменяет незавершённые серии `_process_url` и обнуляет счётчик, чтобы
|
||||
клики по устаревшим ссылкам не «догоняли» уже новый стрим.
|
||||
Возвращает число отброшенных кликов.
|
||||
"""
|
||||
dropped = params.pending_visits
|
||||
if url_tasks:
|
||||
pending = [t for t in url_tasks if not t.done()]
|
||||
for t in pending:
|
||||
t.cancel()
|
||||
if pending:
|
||||
await asyncio.gather(*pending, return_exceptions=True)
|
||||
# Отброшенные клики не будут выполнены — убираем их и из плана,
|
||||
# иначе счётчик "выполнено/запланировано" никогда не сойдётся
|
||||
params.total_planned_visits = max(0, params.total_planned_visits - dropped)
|
||||
params.pending_visits = 0
|
||||
params.remaining_clicks = 0
|
||||
return dropped
|
||||
|
||||
async def _stream_watcher(self, task_id: str, params: TaskParams, message: Message,
|
||||
check_interval: int = 60) -> None:
|
||||
check_interval: int = 60, url_tasks: Optional[set] = None) -> None:
|
||||
"""Проверяет онлайн-статус стрима, ставит/снимает паузу и сохраняет состояние."""
|
||||
was_offline = params.stream_offline
|
||||
|
||||
@@ -2127,6 +2148,10 @@ class BotInterface:
|
||||
|
||||
elif viewers > 0 and was_offline:
|
||||
was_offline = False
|
||||
# Началась новая трансляция — чистим очередь прошлой.
|
||||
# Делаем это до снятия паузы, иначе отложенные серии успеют
|
||||
# проснуться и отработать по неактуальным ссылкам.
|
||||
dropped = await self._drop_pending_series(params, url_tasks)
|
||||
params.stream_offline = False
|
||||
if params.auto_paused:
|
||||
params.paused = False
|
||||
@@ -2134,11 +2159,14 @@ class BotInterface:
|
||||
await self.task_manager.update_task(task_id, paused=False)
|
||||
if self.storage:
|
||||
await self.storage.save_task(task_id, params)
|
||||
logger.info(f"[watcher] {params.channel} is online ({viewers} viewers)")
|
||||
msg = await send_message_safe(
|
||||
message.bot, params.chat_id,
|
||||
f"▶️ {params.channel} снова онлайн ({viewers} зрителей) — мониторинг возобновлён"
|
||||
logger.info(
|
||||
f"[watcher] {params.channel} is online ({viewers} viewers), "
|
||||
f"queue cleared: {dropped} clicks dropped"
|
||||
)
|
||||
text = f"▶️ {params.channel} снова онлайн ({viewers} зрителей) — мониторинг возобновлён"
|
||||
if dropped:
|
||||
text += f"\n🧹 Очередь прошлой трансляции очищена: {dropped} кликов отброшено"
|
||||
msg = await send_message_safe(message.bot, params.chat_id, text)
|
||||
if msg and params.user_id:
|
||||
await self._track_msg(params.user_id, msg.message_id)
|
||||
|
||||
@@ -2170,7 +2198,9 @@ class BotInterface:
|
||||
|
||||
logger.info(f"🚀 GQL chat monitor: {params.channel}")
|
||||
|
||||
watcher = asyncio.create_task(self._stream_watcher(task_id, params, message))
|
||||
watcher = asyncio.create_task(
|
||||
self._stream_watcher(task_id, params, message, url_tasks=url_tasks)
|
||||
)
|
||||
|
||||
try:
|
||||
await irc.listen_for_messages(
|
||||
@@ -2327,6 +2357,7 @@ class BotInterface:
|
||||
|
||||
async def _process_url(self, url: str, username: str, task_id: str, params: TaskParams, message: Message):
|
||||
"""Обрабатывает найденную ссылку (выполняется параллельно)."""
|
||||
notif_msg = None
|
||||
try:
|
||||
params.links_found += 1
|
||||
|
||||
@@ -2438,7 +2469,15 @@ class BotInterface:
|
||||
if params.stopped:
|
||||
return
|
||||
|
||||
# Серия завершена — удаляем уведомление
|
||||
except asyncio.CancelledError:
|
||||
# Очередь сброшена (началась новая трансляция) — серия прекращается
|
||||
logger.info(f"🧹 Серия отменена: {url[:40]}")
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.error(f"Process URL error: {e}")
|
||||
finally:
|
||||
# Уведомление о серии убираем при любом исходе: завершение, отмена,
|
||||
# обрыв прокси, исчерпание баланса
|
||||
if notif_msg:
|
||||
try:
|
||||
await message.bot.delete_message(
|
||||
@@ -2448,9 +2487,6 @@ class BotInterface:
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Process URL error: {e}")
|
||||
|
||||
# =========================================================================
|
||||
# КАБИНЕТ ПОЛЬЗОВАТЕЛЯ
|
||||
# =========================================================================
|
||||
|
||||
Reference in New Issue
Block a user