/
artyuhovs
/
AIFeedBackTrainingBot
Обзор
Документация
Войти
/
artyuhovs
/
AIFeedBackTrainingBot
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
app/bot/progress.py
280 строк
9 KB
Codex
Добавить прозрачный медиаконвейер в версии 0.8.1
19 июл 2026, 20:24
19 июл 2026, 20:24
9c5c8db
Код
Авторство
О чём код?
from __future__ import annotations import asyncio import logging import time from collections.abc import Awaitable, Callable from dataclasses import dataclass from typing import Any from aiogram.exceptions import TelegramAPIError, TelegramRetryAfter logger = logging.getLogger(__name__) CHECK_MARK = "✔️" DEFAULT_MAX_EDITS = 8 DEFAULT_MIN_EDIT_INTERVAL_SECONDS = 1.0 SendProgressMessage = Callable[[str], Awaitable[Any]] Sleep = Callable[[float], Awaitable[None]] Monotonic = Callable[[], float] @dataclass class _ProgressStep: text: str done: bool = False class TelegramProgressReporter: def __init__( self, send_message: SendProgressMessage, *, min_edit_interval_seconds: float = DEFAULT_MIN_EDIT_INTERVAL_SECONDS, max_edits: int = DEFAULT_MAX_EDITS, sleep: Sleep = asyncio.sleep, monotonic: Monotonic = time.monotonic, ) -> None: self._send_message = send_message self._min_edit_interval_seconds = max(min_edit_interval_seconds, 0.0) self._max_edits = max(max_edits, 0) self._sleep = sleep self._monotonic = monotonic self._steps: list[_ProgressStep] = [] self._detail_text: str | None = None self._message: Any | None = None self._service_post_created = False self._last_rendered = "" self._last_edit_at = 0.0 self._edit_count = 0 self._dirty = False self._finalized = False self._pending_task: asyncio.Task[None] | None = None self._lock = asyncio.Lock() async def __call__(self, text: str) -> None: step_text = self._normalize(text) if not step_text: return async with self._lock: if self._finalized: return if not self._service_post_created: self._steps.append(_ProgressStep(step_text)) await self._send_initial_locked() return if self._steps and not self._steps[-1].done and self._steps[-1].text == step_text: return if self._steps: self._steps[-1].done = True self._steps.append(_ProgressStep(step_text)) self._dirty = True await self._maybe_edit_locked(force=False) async def update_current(self, text: str) -> None: """Replace the active status line, for example while a queue position changes.""" step_text = self._normalize(text) if not step_text: return async with self._lock: if self._finalized: return if not self._service_post_created: self._steps.append(_ProgressStep(step_text)) await self._send_initial_locked() return if self._steps and not self._steps[-1].done: self._steps[-1].text = step_text else: self._steps.append(_ProgressStep(step_text)) self._dirty = True # Queue changes are infrequent and must remain visible even after the # ordinary progress edit budget has been exhausted. await self._maybe_edit_locked(force=True, count_edit=False) async def complete( self, final_text: str | None = None, *, detail_text: str | None = None, ) -> None: async with self._lock: if self._finalized: return self._finalized = True if not self._steps: return self._steps[-1].done = True normalized = self._normalize(final_text or "") if normalized: self._steps.append(_ProgressStep(normalized, done=True)) self._detail_text = self._normalize_detail(detail_text) self._dirty = True await self._maybe_edit_locked(force=True) async def fail(self, final_text: str | None = None) -> None: async with self._lock: if self._finalized: return self._finalized = True if not self._steps: return normalized = self._normalize(final_text or "") if normalized: self._steps.append(_ProgressStep(normalized)) self._dirty = True await self._maybe_edit_locked(force=True) async def _send_initial_locked(self) -> None: rendered = self._render() try: self._message = await self._send_message(rendered) except Exception: logger.warning( "telegram.progress_send_failed", extra={"operation": "telegram.progress_send_failed"}, exc_info=True, ) return self._service_post_created = True self._last_rendered = rendered self._last_edit_at = self._monotonic() async def _maybe_edit_locked(self, *, force: bool, count_edit: bool = True) -> None: if not self._service_post_created or self._message is None or not self._dirty: return if self._edit_count >= self._max_edits and not force: return delay = self._edit_delay_locked() if delay > 0: if force: self._cancel_pending_locked() await self._sleep(delay) else: self._schedule_edit_locked(delay) return self._cancel_pending_locked() await self._edit_now_locked(force=force, count_edit=count_edit) async def _edit_now_locked(self, *, force: bool, count_edit: bool) -> None: rendered = self._render() if rendered == self._last_rendered: self._dirty = False return progress_message = self._message if progress_message is None: return try: await progress_message.edit_text(rendered) except TelegramRetryAfter as error: self._dirty = True retry_after = max(float(error.retry_after), 0.0) if force: await self._sleep(retry_after) await self._retry_final_edit_locked(count_edit=count_edit) else: self._schedule_edit_locked(retry_after) return except TelegramAPIError: self._dirty = False logger.warning( "telegram.progress_edit_failed", extra={"operation": "telegram.progress_edit_failed"}, exc_info=True, ) return except Exception: self._dirty = False logger.warning( "telegram.progress_edit_failed", extra={"operation": "telegram.progress_edit_failed"}, exc_info=True, ) return self._mark_edit_sent_locked(rendered, count_edit=count_edit) async def _retry_final_edit_locked(self, *, count_edit: bool) -> None: rendered = self._render() if rendered == self._last_rendered: self._dirty = False return progress_message = self._message if progress_message is None: return try: await progress_message.edit_text(rendered) except Exception: self._dirty = False logger.warning( "telegram.progress_retry_failed", extra={"operation": "telegram.progress_retry_failed"}, exc_info=True, ) return self._mark_edit_sent_locked(rendered, count_edit=count_edit) def _mark_edit_sent_locked(self, rendered: str, *, count_edit: bool) -> None: self._last_rendered = rendered self._last_edit_at = self._monotonic() if count_edit: self._edit_count += 1 self._dirty = False def _edit_delay_locked(self) -> float: if self._last_edit_at <= 0: return 0.0 elapsed = self._monotonic() - self._last_edit_at return max(self._min_edit_interval_seconds - elapsed, 0.0) def _schedule_edit_locked(self, delay: float) -> None: if self._pending_task is not None and not self._pending_task.done(): return self._pending_task = asyncio.create_task(self._delayed_edit(delay)) def _cancel_pending_locked(self) -> None: if self._pending_task is not None and not self._pending_task.done(): self._pending_task.cancel() self._pending_task = None async def _delayed_edit(self, delay: float) -> None: try: await self._sleep(delay) except asyncio.CancelledError: return async with self._lock: self._pending_task = None await self._maybe_edit_locked(force=False) def _render(self) -> str: steps = "\n".join(self._render_step(step) for step in self._steps) if self._detail_text: return f"{steps}\n\n{self._detail_text}" return steps @staticmethod def _render_step(step: _ProgressStep) -> str: prefix = f"{CHECK_MARK} " if step.done else "" return f"{prefix}{step.text}" @staticmethod def _normalize(text: str) -> str: return " ".join(text.strip().split()) @staticmethod def _normalize_detail(text: str | None) -> str | None: normalized = (text or "").strip() return normalized or None