FreshPixels commited on
Commit
7e686b5
·
verified ·
1 Parent(s): 63c9475

Upload 9 files

Browse files
Files changed (9) hide show
  1. bot.py +204 -0
  2. config.py +129 -0
  3. database.py +403 -0
  4. handlers callbacks.py +82 -0
  5. handlers chat.py +261 -0
  6. handlers commands.py +149 -0
  7. middlewares rate_limit.py +43 -0
  8. services glm.py +344 -0
  9. utils helpers.py +114 -0
bot.py ADDED
@@ -0,0 +1,204 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import asyncio
2
+ import logging
3
+ import os
4
+ import ssl
5
+ import sys
6
+
7
+ import aiohttp
8
+ import certifi
9
+ from aiohttp import web
10
+ from aiogram import Bot, Dispatcher
11
+ from aiogram.enums import ParseMode
12
+ from aiogram.client.default import DefaultBotProperties
13
+ from aiogram.webhook.aiohttp_server import SimpleRequestHandler, setup_application
14
+
15
+ from config import config
16
+ from database import db
17
+ from handlers import commands_router, chat_router, callbacks_router
18
+ from middlewares import OwnerMiddleware, RateLimitMiddleware
19
+
20
+
21
+ # ═══════════════════════════════════════════════════════════════════
22
+ # BLOCK 0: SSL-патчи для HF Spaces (certifi + proxy)
23
+ # ═══════════════════════════════════════════════════════════════════
24
+ custom_ssl = ssl.create_default_context(cafile=certifi.where())
25
+ custom_ssl.check_hostname = True
26
+ custom_ssl.verify_mode = ssl.CERT_REQUIRED
27
+
28
+ _orig_tcp_init = aiohttp.TCPConnector.__init__
29
+
30
+
31
+ def _patched_tcp_init(self, *args, **kwargs):
32
+ if kwargs.get("ssl") is not False:
33
+ kwargs["ssl"] = custom_ssl
34
+ _orig_tcp_init(self, *args, **kwargs)
35
+
36
+
37
+ aiohttp.TCPConnector.__init__ = _patched_tcp_init
38
+
39
+ _orig_session_init = aiohttp.ClientSession.__init__
40
+
41
+
42
+ def _patched_session_init(self, *args, **kwargs):
43
+ kwargs["trust_env"] = True
44
+ _orig_session_init(self, *args, **kwargs)
45
+
46
+
47
+ aiohttp.ClientSession.__init__ = _patched_session_init
48
+
49
+ _orig_request = aiohttp.ClientSession._request
50
+
51
+
52
+ async def _patched_request(self, method, url, *args, **kwargs):
53
+ proxy_server = os.getenv("TELEGRAM_API_SERVER")
54
+ if proxy_server and "api.telegram.org" in str(url):
55
+ proxy_server = proxy_server.strip().rstrip("/")
56
+ str_url = str(url).replace("https://api.telegram.org", proxy_server)
57
+ logging.info("🔀 Переадресация aiogram через прокси ➡️ %s", str_url)
58
+ url = str_url
59
+ return await _orig_request(self, method, url, *args, **kwargs)
60
+
61
+
62
+ aiohttp.ClientSession._request = _patched_request
63
+
64
+
65
+ # ═══════════════════════════════════════════════════════════════════
66
+ # BLOCK 1: Логирование
67
+ # ═══════════════════════════════════════════════════════════════════
68
+ def setup_logging() -> None:
69
+ level = getattr(logging, config.LOG_LEVEL.upper(), logging.INFO)
70
+ logging.basicConfig(
71
+ level=level,
72
+ format="%(asctime)s - %(name)s - %(levelname)s - %(message)s",
73
+ stream=sys.stdout,
74
+ )
75
+ # Уменьшаем шум от библиотек
76
+ logging.getLogger("httpx").setLevel(logging.WARNING)
77
+ logging.getLogger("httpcore").setLevel(logging.WARNING)
78
+ logging.getLogger("aiogram").setLevel(logging.INFO)
79
+
80
+
81
+ # ═══════════════════════════════════════════════════════════════════
82
+ # BLOCK 2: Инициализация бота и диспетчера
83
+ # ═══════════════════════════════════════════════════════════════════
84
+ bot = Bot(
85
+ token=config.BOT_TOKEN,
86
+ default=DefaultBotProperties(parse_mode=ParseMode.HTML)
87
+ )
88
+ dp = Dispatcher()
89
+
90
+ # Middlewares (порядок важен!)
91
+ dp.message.middleware(OwnerMiddleware())
92
+ dp.message.middleware(RateLimitMiddleware())
93
+
94
+ # Logging middleware
95
+ @dp.update.middleware()
96
+ async def log_updates(handler, event, data):
97
+ logger = logging.getLogger(__name__)
98
+ logger.info("→ Update received: %s", event.update_id if hasattr(event, 'update_id') else 'N/A')
99
+ try:
100
+ return await handler(event, data)
101
+ except Exception as e:
102
+ logger.error("Update failed: %s", e, exc_info=True)
103
+ raise
104
+
105
+ # Routers (порядок важен: commands → callbacks → chat)
106
+ dp.include_router(commands_router)
107
+ dp.include_router(callbacks_router)
108
+ dp.include_router(chat_router)
109
+
110
+
111
+ # ═══════════════════════════════════════════════════════════════════
112
+ # BLOCK 3: HTTP Handlers
113
+ # ═══════════════════════════════════════════════════════════════════
114
+ @web.middleware
115
+ async def hf_logging_middleware(request, handler):
116
+ return await handler(request)
117
+
118
+
119
+ async def health_check(request: web.Request) -> web.Response:
120
+ """HF Spaces Health Check — обязательно отвечать 200 на '/'."""
121
+ return web.Response(text="🚀 GLM Bot РАБОТАЕТ!")
122
+
123
+
124
+ # ═══════════════════════════════════════════════════════════════════
125
+ # BLOCK 4: Lifecycle hooks
126
+ # ═══════════════════════════════════════════════════════════════════
127
+ async def on_startup(app: web.Application) -> None:
128
+ logger = logging.getLogger(__name__)
129
+ await db.connect()
130
+ logger.info("✅ Database connected")
131
+
132
+ space_host = os.getenv("SPACE_HOST", "")
133
+ if space_host:
134
+ full_webhook_link = f"https://{space_host.strip()}{config.WEBHOOK_PATH}"
135
+ for attempt in range(5):
136
+ try:
137
+ await bot.set_webhook(
138
+ url=full_webhook_link,
139
+ drop_pending_updates=True,
140
+ request_timeout=30,
141
+ )
142
+ logger.info("✅ Webhook установлен: %s", full_webhook_link)
143
+ break
144
+ except Exception as e:
145
+ logger.warning(
146
+ "⚠️ Попытка %d/5 установки webhook: %s", attempt + 1, e
147
+ )
148
+ await asyncio.sleep(5)
149
+ else:
150
+ logger.warning("⚠️ SPACE_HOST не задан, webhook не установлен!")
151
+
152
+
153
+ async def on_shutdown(app: web.Application) -> None:
154
+ logger = logging.getLogger(__name__)
155
+ logger.info("🛑 Shutdown начат...")
156
+
157
+ try:
158
+ await bot.delete_webhook(drop_pending_updates=True)
159
+ logger.info("Webhook удалён")
160
+ except Exception as e:
161
+ logger.warning("Ошибка при удалении webhook: %s", e)
162
+
163
+ try:
164
+ await dp.storage.close()
165
+ await bot.session.close()
166
+ logger.info("Bot session закрыт")
167
+ except Exception as e:
168
+ logger.warning("Ошибка при закрытии сессии бота: %s", e)
169
+
170
+ try:
171
+ await db.disconnect()
172
+ logger.info("Database disconnected")
173
+ except Exception as e:
174
+ logger.warning("Ошибка при отключении БД: %s", e)
175
+
176
+
177
+ # ═══════════════════════════════════════════════════════════════════
178
+ # BLOCK 5: Main
179
+ # ═══════════════════════════════════════════════════════════════════
180
+ def main() -> None:
181
+ setup_logging()
182
+ logger = logging.getLogger(__name__)
183
+ logger.info("🚀 Starting HF Spaces bot...")
184
+ logger.info("Config: model=%s, fallback=%s, streaming=%s, rate_limit=%s",
185
+ config.PRIMARY_MODEL, config.FALLBACK_MODEL,
186
+ config.STREAMING_ENABLED, config.RATE_LIMIT_ENABLED)
187
+
188
+ app = web.Application(middlewares=[hf_logging_middleware])
189
+ app.router.add_get("/", health_check)
190
+
191
+ webhook_requests_handler = SimpleRequestHandler(dispatcher=dp, bot=bot)
192
+ webhook_requests_handler.register(app, path=config.WEBHOOK_PATH)
193
+ setup_application(app, dp, bot=bot)
194
+
195
+ app.on_startup.append(on_startup)
196
+ app.on_shutdown.append(on_shutdown)
197
+
198
+ port = int(os.environ.get("PORT", 7860))
199
+ logger.info("🚀 Запуск сервера на %s:%d...", "0.0.0.0", port)
200
+ web.run_app(app, host="0.0.0.0", port=port)
201
+
202
+
203
+ if __name__ == "__main__":
204
+ main()
config.py ADDED
@@ -0,0 +1,129 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import os
2
+ from dataclasses import dataclass, field
3
+ from dotenv import load_dotenv
4
+
5
+ load_dotenv()
6
+
7
+
8
+ @dataclass
9
+ class Config:
10
+ # ═══════════════════════════════════════════════════════════════
11
+ # BLOCK: Telegram
12
+ # ═══════════════════════════════════════════════════════════════
13
+ BOT_TOKEN: str = os.getenv("BOT_TOKEN", "")
14
+ OWNER_ID: int = int(os.getenv("OWNER_ID", "0") or "0")
15
+ WEBHOOK_PATH: str = "/webhook"
16
+
17
+ # ═══════════════════════════════════════════════════════════════
18
+ # BLOCK: Database (Neon PostgreSQL)
19
+ # ═══════════════════════════════════════════════════════════════
20
+ DATABASE_URL: str = os.getenv("DATABASE_URL", "")
21
+
22
+ # ═══════════════════════════════════════════════════════════════
23
+ # BLOCK: AI — Primary Model (GLM-4.5 via NVIDIA)
24
+ # ═══════════════════════════════════════════════════════════════
25
+ NVIDIA_API_KEY: str = os.getenv("NVIDIA_API_KEY", "")
26
+ NVIDIA_BASE_URL: str = "https://integrate.api.nvidia.com/v1"
27
+ PRIMARY_MODEL: str = "z-ai/glm-4.5" # или "z-ai/glm-5.1" — обнови под свою модель
28
+
29
+ # ═══════════════════════════════════════════════════════════════
30
+ # BLOCK: AI — Fallback Model
31
+ # ═══════════════════════════════════════════════════════════════
32
+ FALLBACK_MODEL: str = os.getenv("FALLBACK_MODEL", "meta/llama-3.3-70b-instruct")
33
+ FALLBACK_ENABLED: bool = os.getenv("FALLBACK_ENABLED", "true").lower() == "true"
34
+
35
+ # ═══════════════════════════════════════════════════════════════
36
+ # BLOCK: AI Parameters — Optimized for GLM-4.5
37
+ # ═══════════════════════════════════════════════════════════════
38
+ # GLM-4.5 рекомендации через NVIDIA NIM:
39
+ # - temperature: 0.3–0.8 (ниже = точнее, выше = креативнее)
40
+ # - top_p: 0.7–0.95
41
+ # - max_tokens: зависит от задачи, 4096–8192 для длинных ответов
42
+ GLM_TEMPERATURE: float = float(os.getenv("GLM_TEMPERATURE", "0.7"))
43
+ GLM_TOP_P: float = float(os.getenv("GLM_TOP_P", "0.9"))
44
+ GLM_FREQUENCY_PENALTY: float = float(os.getenv("GLM_FREQUENCY_PENALTY", "0.1"))
45
+ GLM_PRESENCE_PENALTY: float = float(os.getenv("GLM_PRESENCE_PENALTY", "0.05"))
46
+ GLM_MAX_TOKENS: int = int(os.getenv("GLM_MAX_TOKENS", "4096"))
47
+
48
+ # ═══════════════════════════════════════════════════════════════
49
+ # BLOCK: Timeouts — Granular for NVIDIA NIM
50
+ # ═══════════════════════════════════════════════════════════════
51
+ # NVIDIA NIM может отвечать 2–3 минуты на сложные запросы.
52
+ # connect: время установления TCP-соединения
53
+ # read: время ожидания первого байта ответа
54
+ # write: время на отправку запроса
55
+ TIMEOUT_CONNECT: float = float(os.getenv("TIMEOUT_CONNECT", "10.0"))
56
+ TIMEOUT_READ: float = float(os.getenv("TIMEOUT_READ", "180.0")) # 3 минуты
57
+ TIMEOUT_WRITE: float = float(os.getenv("TIMEOUT_WRITE", "10.0"))
58
+ TIMEOUT_POOL: float = float(os.getenv("TIMEOUT_POOL", "5.0"))
59
+
60
+ # ═══════════════════════════════════════════════════════════════
61
+ # BLOCK: Retry Configuration
62
+ # ═══════════════════════════════════════════════════════════════
63
+ MAX_RETRIES: int = int(os.getenv("MAX_RETRIES", "3"))
64
+ RETRY_BASE_DELAY: float = float(os.getenv("RETRY_BASE_DELAY", "2.0"))
65
+ RETRY_MAX_DELAY: float = float(os.getenv("RETRY_MAX_DELAY", "60.0"))
66
+ RETRY_EXPONENTIAL_BASE: float = float(os.getenv("RETRY_EXPONENTIAL_BASE", "2.0"))
67
+
68
+ # ═══════════════════════════════════════════════════════════════
69
+ # BLOCK: Streaming
70
+ # ═══════════════════════════════════════════════════════════════
71
+ STREAMING_ENABLED: bool = os.getenv("STREAMING_ENABLED", "true").lower() == "true"
72
+ STREAMING_CHUNK_SIZE: int = int(os.getenv("STREAMING_CHUNK_SIZE", "100"))
73
+ STREAMING_UPDATE_INTERVAL: float = float(os.getenv("STREAMING_UPDATE_INTERVAL", "1.5"))
74
+
75
+ # ═══════════════════════════════════════════════════════════════
76
+ # BLOCK: History & Memory
77
+ # ═══════════════════════════════════════════════════════════════
78
+ MAX_HISTORY: int = int(os.getenv("MAX_HISTORY", "20")) # сообщений в контексте
79
+ MAX_CONTEXT_TOKENS: int = int(os.getenv("MAX_CONTEXT_TOKENS", "6000")) # примерно токенов
80
+ SUMMARIZE_THRESHOLD: int = int(os.getenv("SUMMARIZE_THRESHOLD", "30"))
81
+ SUMMARY_MAX_TOKENS: int = int(os.getenv("SUMMARY_MAX_TOKENS", "512"))
82
+ KEEP_RECENT_MESSAGES: int = int(os.getenv("KEEP_RECENT_MESSAGES", "10"))
83
+
84
+ # ═══════════════════════════════════════════════════════════════
85
+ # BLOCK: Telegram Limits
86
+ # ═══════════════════════════════════════════════════════════════
87
+ MAX_MESSAGE_LENGTH: int = 4096
88
+ TYPING_ACTION_INTERVAL: float = 4.5 # Telegram typing action expires after ~5s
89
+
90
+ # ═══════════════════════════════════════════════════════════════
91
+ # BLOCK: Rate Limiting
92
+ # ═══════════════════════════════════════════════════════════════
93
+ RATE_LIMIT_ENABLED: bool = os.getenv("RATE_LIMIT_ENABLED", "true").lower() == "true"
94
+ RATE_LIMIT_REQUESTS_PER_MINUTE: int = int(os.getenv("RATE_LIMIT_REQUESTS_PER_MINUTE", "10"))
95
+ RATE_LIMIT_BURST: int = int(os.getenv("RATE_LIMIT_BURST", "3"))
96
+
97
+ # ═══════════════════════════════════════════════════════════════
98
+ # BLOCK: Cache
99
+ # ═══════════════════════════════════════════════════════════════
100
+ CACHE_ENABLED: bool = os.getenv("CACHE_ENABLED", "true").lower() == "true"
101
+ CACHE_TTL_SECONDS: int = int(os.getenv("CACHE_TTL_SECONDS", "300")) # 5 минут
102
+ CACHE_MAX_SIZE: int = int(os.getenv("CACHE_MAX_SIZE", "1000"))
103
+
104
+ # ═══════════════════════════════════════════════════════════════
105
+ # BLOCK: Monitoring
106
+ # ═══════════════════════════════════════════════════════════════
107
+ LOG_LEVEL: str = os.getenv("LOG_LEVEL", "INFO")
108
+ METRICS_ENABLED: bool = os.getenv("METRICS_ENABLED", "true").lower() == "true"
109
+
110
+ def __post_init__(self) -> None:
111
+ if self.DATABASE_URL.startswith("postgres://"):
112
+ self.DATABASE_URL = self.DATABASE_URL.replace(
113
+ "postgres://", "postgresql://", 1
114
+ )
115
+
116
+ def validate(self) -> None:
117
+ required = [
118
+ ("BOT_TOKEN", self.BOT_TOKEN),
119
+ ("OWNER_ID", self.OWNER_ID),
120
+ ("DATABASE_URL", self.DATABASE_URL),
121
+ ("NVIDIA_API_KEY", self.NVIDIA_API_KEY),
122
+ ]
123
+ for name, value in required:
124
+ if not value:
125
+ raise ValueError(f"{name} is required")
126
+
127
+
128
+ config = Config()
129
+ config.validate()
database.py ADDED
@@ -0,0 +1,403 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import asyncpg
2
+ import logging
3
+ import time
4
+ from typing import Optional, List, Dict, Any
5
+ from config import config
6
+
7
+ logger = logging.getLogger(__name__)
8
+
9
+
10
+ class Database:
11
+ def __init__(self) -> None:
12
+ self.pool: Optional[asyncpg.Pool] = None
13
+
14
+ async def connect(self) -> None:
15
+ self.pool = await asyncpg.create_pool(
16
+ dsn=config.DATABASE_URL,
17
+ min_size=1,
18
+ max_size=5, # увеличено для масштабируемости
19
+ command_timeout=60,
20
+ server_settings={
21
+ "jit": "off", # отключаем JIT для стабильности
22
+ "application_name": "glm_bot",
23
+ },
24
+ )
25
+ logger.info("Database pool created (max_size=5)")
26
+ await self._create_tables()
27
+ await self._create_indexes()
28
+
29
+ async def disconnect(self) -> None:
30
+ if self.pool:
31
+ await self.pool.close()
32
+ logger.info("Database pool closed")
33
+
34
+ def _acquire(self):
35
+ if self.pool is None:
36
+ raise RuntimeError("Database not connected. Call connect() first.")
37
+ return self.pool.acquire()
38
+
39
+ async def _create_tables(self) -> None:
40
+ async with self._acquire() as conn:
41
+ # Users table
42
+ await conn.execute("""
43
+ CREATE TABLE IF NOT EXISTS users (
44
+ id BIGINT PRIMARY KEY,
45
+ username VARCHAR(255),
46
+ first_name VARCHAR(255),
47
+ last_name VARCHAR(255),
48
+ created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
49
+ updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
50
+ settings JSONB DEFAULT '{}'::jsonb
51
+ )
52
+ """)
53
+
54
+ # Messages table
55
+ await conn.execute("""
56
+ CREATE TABLE IF NOT EXISTS messages (
57
+ id SERIAL PRIMARY KEY,
58
+ user_id BIGINT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
59
+ role VARCHAR(20) NOT NULL CHECK (role IN ('user', 'assistant', 'system')),
60
+ content TEXT NOT NULL,
61
+ tokens_used INTEGER DEFAULT 0,
62
+ is_summarized BOOLEAN DEFAULT FALSE,
63
+ created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
64
+ )
65
+ """)
66
+
67
+ # Summaries table
68
+ await conn.execute("""
69
+ CREATE TABLE IF NOT EXISTS summaries (
70
+ id SERIAL PRIMARY KEY,
71
+ user_id BIGINT NOT NULL UNIQUE REFERENCES users(id) ON DELETE CASCADE,
72
+ summary TEXT NOT NULL,
73
+ message_count INTEGER NOT NULL DEFAULT 0,
74
+ created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
75
+ updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
76
+ )
77
+ """)
78
+
79
+ # Metrics table — для мониторинга
80
+ await conn.execute("""
81
+ CREATE TABLE IF NOT EXISTS metrics (
82
+ id SERIAL PRIMARY KEY,
83
+ user_id BIGINT REFERENCES users(id) ON DELETE SET NULL,
84
+ model VARCHAR(100),
85
+ request_duration_ms FLOAT,
86
+ tokens_input INTEGER DEFAULT 0,
87
+ tokens_output INTEGER DEFAULT 0,
88
+ total_tokens INTEGER DEFAULT 0,
89
+ success BOOLEAN DEFAULT TRUE,
90
+ error_type VARCHAR(100),
91
+ created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
92
+ )
93
+ """)
94
+
95
+ # Rate limiting table
96
+ await conn.execute("""
97
+ CREATE TABLE IF NOT EXISTS rate_limits (
98
+ user_id BIGINT PRIMARY KEY REFERENCES users(id) ON DELETE CASCADE,
99
+ request_count INTEGER DEFAULT 0,
100
+ window_start TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
101
+ updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
102
+ )
103
+ """)
104
+
105
+ logger.info("Database tables created/verified")
106
+
107
+ async def _create_indexes(self) -> None:
108
+ async with self._acquire() as conn:
109
+ await conn.execute("""
110
+ CREATE INDEX IF NOT EXISTS idx_messages_user_id_created_at
111
+ ON messages(user_id, created_at DESC)
112
+ """)
113
+ await conn.execute("""
114
+ CREATE INDEX IF NOT EXISTS idx_messages_user_id_summarized
115
+ ON messages(user_id, is_summarized, created_at DESC)
116
+ """)
117
+ await conn.execute("""
118
+ CREATE INDEX IF NOT EXISTS idx_metrics_user_created
119
+ ON metrics(user_id, created_at DESC)
120
+ """)
121
+ await conn.execute("""
122
+ CREATE INDEX IF NOT EXISTS idx_metrics_created_at
123
+ ON metrics(created_at DESC)
124
+ """)
125
+
126
+ # ═══════════════════════════════════════════════════════════════
127
+ # Users
128
+ # ═══════════════════════════════════════════════════════════════
129
+ async def upsert_user(
130
+ self,
131
+ user_id: int,
132
+ username: Optional[str],
133
+ first_name: Optional[str],
134
+ last_name: Optional[str],
135
+ ) -> None:
136
+ async with self._acquire() as conn:
137
+ await conn.execute("""
138
+ INSERT INTO users (id, username, first_name, last_name)
139
+ VALUES ($1, $2, $3, $4)
140
+ ON CONFLICT (id) DO UPDATE SET
141
+ username = EXCLUDED.username,
142
+ first_name = EXCLUDED.first_name,
143
+ last_name = EXCLUDED.last_name,
144
+ updated_at = NOW()
145
+ """, user_id, username, first_name, last_name)
146
+
147
+ async def get_user_settings(self, user_id: int) -> Dict[str, Any]:
148
+ async with self._acquire() as conn:
149
+ row = await conn.fetchrow("""
150
+ SELECT settings FROM users WHERE id = $1
151
+ """, user_id)
152
+ return row["settings"] if row and row["settings"] else {}
153
+
154
+ async def update_user_settings(self, user_id: int, settings: Dict[str, Any]) -> None:
155
+ async with self._acquire() as conn:
156
+ await conn.execute("""
157
+ UPDATE users SET settings = $2, updated_at = NOW() WHERE id = $1
158
+ """, user_id, settings)
159
+
160
+ # ═══════════════════════════════════════════════════════════════
161
+ # Messages
162
+ # ═══════════════════════════════════════════════════════════════
163
+ async def save_message(
164
+ self, user_id: int, role: str, content: str, tokens_used: int = 0
165
+ ) -> None:
166
+ async with self._acquire() as conn:
167
+ await conn.execute("""
168
+ INSERT INTO messages (user_id, role, content, tokens_used)
169
+ VALUES ($1, $2, $3, $4)
170
+ """, user_id, role, content, tokens_used)
171
+
172
+ async def get_messages(self, user_id: int, limit: int = 30) -> List[Dict[str, Any]]:
173
+ async with self._acquire() as conn:
174
+ rows = await conn.fetch("""
175
+ SELECT id, role, content, tokens_used, created_at
176
+ FROM messages
177
+ WHERE user_id = $1 AND is_summarized = FALSE
178
+ ORDER BY created_at DESC
179
+ LIMIT $2
180
+ """, user_id, limit)
181
+ return [
182
+ {
183
+ "id": r["id"],
184
+ "role": r["role"],
185
+ "content": r["content"],
186
+ "tokens_used": r["tokens_used"],
187
+ "created_at": r["created_at"],
188
+ }
189
+ for r in reversed(rows)
190
+ ]
191
+
192
+ async def get_messages_with_token_budget(
193
+ self, user_id: int, max_tokens: int
194
+ ) -> List[Dict[str, Any]]:
195
+ """Возвращает сообщения, укладывающиеся в бюджет токенов."""
196
+ async with self._acquire() as conn:
197
+ rows = await conn.fetch("""
198
+ SELECT id, role, content, tokens_used, created_at
199
+ FROM messages
200
+ WHERE user_id = $1 AND is_summarized = FALSE
201
+ ORDER BY created_at DESC
202
+ """, user_id)
203
+
204
+ result = []
205
+ total_tokens = 0
206
+ for r in rows:
207
+ msg_tokens = r["tokens_used"] or len(r["content"].split()) * 2
208
+ if total_tokens + msg_tokens > max_tokens and result:
209
+ break
210
+ total_tokens += msg_tokens
211
+ result.insert(0, {
212
+ "id": r["id"],
213
+ "role": r["role"],
214
+ "content": r["content"],
215
+ "tokens_used": r["tokens_used"],
216
+ "created_at": r["created_at"],
217
+ })
218
+ return result
219
+
220
+ # ═══════════════════════════════════════════════════════════════
221
+ # Summaries
222
+ # ═══════════════════════════════════════════════════════════════
223
+ async def get_summary(self, user_id: int) -> Optional[str]:
224
+ async with self._acquire() as conn:
225
+ row = await conn.fetchrow("""
226
+ SELECT summary FROM summaries WHERE user_id = $1
227
+ """, user_id)
228
+ return row["summary"] if row else None
229
+
230
+ async def save_summary(self, user_id: int, summary: str, message_count: int) -> None:
231
+ async with self._acquire() as conn:
232
+ await conn.execute("""
233
+ INSERT INTO summaries (user_id, summary, message_count, updated_at)
234
+ VALUES ($1, $2, $3, NOW())
235
+ ON CONFLICT (user_id) DO UPDATE SET
236
+ summary = EXCLUDED.summary,
237
+ message_count = summaries.message_count + EXCLUDED.message_count,
238
+ updated_at = NOW()
239
+ """, user_id, summary, message_count)
240
+
241
+ async def mark_summarized(self, user_id: int, cutoff_id: int) -> None:
242
+ async with self._acquire() as conn:
243
+ await conn.execute("""
244
+ UPDATE messages
245
+ SET is_summarized = TRUE
246
+ WHERE user_id = $1 AND id <= $2
247
+ """, user_id, cutoff_id)
248
+
249
+ async def get_oldest_unsummarized(self, user_id: int, limit: int) -> List[Dict[str, Any]]:
250
+ async with self._acquire() as conn:
251
+ rows = await conn.fetch("""
252
+ SELECT id, role, content
253
+ FROM messages
254
+ WHERE user_id = $1 AND is_summarized = FALSE
255
+ ORDER BY created_at ASC
256
+ LIMIT $2
257
+ """, user_id, limit)
258
+ return [{"id": r["id"], "role": r["role"], "content": r["content"]} for r in rows]
259
+
260
+ async def count_unsummarized(self, user_id: int) -> int:
261
+ async with self._acquire() as conn:
262
+ val = await conn.fetchval("""
263
+ SELECT COUNT(*) FROM messages
264
+ WHERE user_id = $1 AND is_summarized = FALSE
265
+ """, user_id)
266
+ return val or 0
267
+
268
+ # ═══════════════════════════════════════════════════════════════
269
+ # Clear History
270
+ # ═══════════════════════════════════════════════════════════════
271
+ async def clear_history(self, user_id: int) -> int:
272
+ async with self._acquire() as conn:
273
+ async with conn.transaction():
274
+ result = await conn.execute("""
275
+ DELETE FROM messages WHERE user_id = $1
276
+ """, user_id)
277
+ await conn.execute("""
278
+ DELETE FROM summaries WHERE user_id = $1
279
+ """, user_id)
280
+ try:
281
+ count = int(result.split()[-1])
282
+ except (ValueError, IndexError):
283
+ count = 0
284
+ logger.info("Cleared %d messages and summary for user %s", count, user_id)
285
+ return count
286
+
287
+ # ═══════════════════════════════════════════════════════════════
288
+ # Stats
289
+ # ═══════════════════════════════════════════════════════════════
290
+ async def get_stats(self, user_id: int) -> Dict[str, Any]:
291
+ async with self._acquire() as conn:
292
+ user_count = await conn.fetchval("SELECT COUNT(*) FROM users")
293
+ msg_count = await conn.fetchval(
294
+ "SELECT COUNT(*) FROM messages WHERE user_id = $1", user_id
295
+ )
296
+ total_msg_count = await conn.fetchval("SELECT COUNT(*) FROM messages")
297
+ summary = await conn.fetchval(
298
+ "SELECT message_count FROM summaries WHERE user_id = $1", user_id
299
+ )
300
+ total_tokens = await conn.fetchval("""
301
+ SELECT COALESCE(SUM(total_tokens), 0) FROM metrics WHERE user_id = $1
302
+ """, user_id)
303
+ avg_latency = await conn.fetchval("""
304
+ SELECT COALESCE(AVG(request_duration_ms), 0)
305
+ FROM metrics WHERE user_id = $1 AND success = TRUE
306
+ """, user_id)
307
+ return {
308
+ "total_users": user_count,
309
+ "user_messages": msg_count,
310
+ "total_messages": total_msg_count,
311
+ "summarized_messages": summary or 0,
312
+ "total_tokens_used": int(total_tokens),
313
+ "avg_latency_ms": round(avg_latency, 1) if avg_latency else 0,
314
+ }
315
+
316
+ # ═══════════════════════════════════════════════════════════════
317
+ # Metrics
318
+ # ═══════════════════════════════════════════════════════════════
319
+ async def save_metric(
320
+ self,
321
+ user_id: Optional[int],
322
+ model: str,
323
+ duration_ms: float,
324
+ tokens_input: int = 0,
325
+ tokens_output: int = 0,
326
+ success: bool = True,
327
+ error_type: Optional[str] = None,
328
+ ) -> None:
329
+ async with self._acquire() as conn:
330
+ await conn.execute("""
331
+ INSERT INTO metrics
332
+ (user_id, model, request_duration_ms, tokens_input, tokens_output,
333
+ total_tokens, success, error_type)
334
+ VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
335
+ """, user_id, model, duration_ms, tokens_input, tokens_output,
336
+ tokens_input + tokens_output, success, error_type)
337
+
338
+ async def get_user_metrics(self, user_id: int, limit: int = 50) -> List[Dict[str, Any]]:
339
+ async with self._acquire() as conn:
340
+ rows = await conn.fetch("""
341
+ SELECT model, request_duration_ms, total_tokens, success, error_type, created_at
342
+ FROM metrics WHERE user_id = $1 ORDER BY created_at DESC LIMIT $2
343
+ """, user_id, limit)
344
+ return [dict(r) for r in rows]
345
+
346
+ # ═══════════════════════════════════════════════════════════════
347
+ # Rate Limiting
348
+ # ═══════════════════════════════════════════════════════════════
349
+ async def check_rate_limit(self, user_id: int) -> tuple[bool, int, float]:
350
+ """Returns (allowed, remaining_requests, reset_in_seconds)."""
351
+ if not config.RATE_LIMIT_ENABLED:
352
+ return True, 999, 0.0
353
+
354
+ async with self._acquire() as conn:
355
+ async with conn.transaction():
356
+ row = await conn.fetchrow("""
357
+ SELECT request_count, window_start
358
+ FROM rate_limits WHERE user_id = $1
359
+ FOR UPDATE
360
+ """, user_id)
361
+
362
+ now = time.time()
363
+ window_duration = 60 # 1 minute
364
+
365
+ if not row:
366
+ await conn.execute("""
367
+ INSERT INTO rate_limits (user_id, request_count, window_start)
368
+ VALUES ($1, 1, NOW())
369
+ """, user_id)
370
+ return True, config.RATE_LIMIT_REQUESTS_PER_MINUTE - 1, window_duration
371
+
372
+ window_start_ts = row["window_start"].timestamp()
373
+ if now - window_start_ts >= window_duration:
374
+ # Window expired, reset
375
+ await conn.execute("""
376
+ UPDATE rate_limits
377
+ SET request_count = 1, window_start = NOW(), updated_at = NOW()
378
+ WHERE user_id = $1
379
+ """, user_id)
380
+ return True, config.RATE_LIMIT_REQUESTS_PER_MINUTE - 1, window_duration
381
+
382
+ if row["request_count"] >= config.RATE_LIMIT_REQUESTS_PER_MINUTE:
383
+ reset_in = window_duration - (now - window_start_ts)
384
+ return False, 0, reset_in
385
+
386
+ await conn.execute("""
387
+ UPDATE rate_limits
388
+ SET request_count = request_count + 1, updated_at = NOW()
389
+ WHERE user_id = $1
390
+ """, user_id)
391
+
392
+ remaining = config.RATE_LIMIT_REQUESTS_PER_MINUTE - row["request_count"] - 1
393
+ reset_in = window_duration - (now - window_start_ts)
394
+ return True, remaining, reset_in
395
+
396
+ async def reset_rate_limit(self, user_id: int) -> None:
397
+ async with self._acquire() as conn:
398
+ await conn.execute("""
399
+ DELETE FROM rate_limits WHERE user_id = $1
400
+ """, user_id)
401
+
402
+
403
+ db = Database()
handlers callbacks.py ADDED
@@ -0,0 +1,82 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import logging
2
+ from aiogram import Router, types, F
3
+ from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton
4
+ from database import db
5
+ from handlers.commands import get_main_menu
6
+
7
+ logger = logging.getLogger(__name__)
8
+ router = Router()
9
+
10
+
11
+ @router.callback_query(F.data == "regenerate")
12
+ async def cb_regenerate(callback: types.CallbackQuery) -> None:
13
+ """Помечает последнее сообщение пользователя для регенерации."""
14
+ await callback.answer("🔄 Регенерация...")
15
+ # Удаляем последний ответ ассистента, чтобы chat handler перегенерировал
16
+ user_id = callback.from_user.id
17
+ # Получаем последнее сообщение пользователя
18
+ history = await db.get_messages(user_id, limit=2)
19
+ if len(history) >= 1 and history[-1]["role"] == "user":
20
+ await callback.message.answer(
21
+ "🔄 <b>Регенерирую ответ...</b>",
22
+ parse_mode="HTML",
23
+ )
24
+ else:
25
+ await callback.message.answer(
26
+ "❌ Нет сообщения для регенерации.",
27
+ parse_mode="HTML",
28
+ )
29
+
30
+
31
+ @router.callback_query(F.data == "clear_memory")
32
+ async def cb_clear_memory(callback: types.CallbackQuery) -> None:
33
+ user_id = callback.from_user.id
34
+ count = await db.clear_history(user_id)
35
+ await callback.answer("🗑 Память очищена!")
36
+ await callback.message.answer(
37
+ f"🗑 <b>Память очищена.</b>\nУдалено сообщений: <code>{count}</code>.",
38
+ parse_mode="HTML",
39
+ reply_markup=get_main_menu(),
40
+ )
41
+
42
+
43
+ @router.callback_query(F.data == "show_stats")
44
+ async def cb_show_stats(callback: types.CallbackQuery) -> None:
45
+ stats = await db.get_stats(callback.from_user.id)
46
+ await callback.answer()
47
+ await callback.message.answer(
48
+ f"📊 <b>Статистика:</b>\n\n"
49
+ f"💬 Ваших сообщений: <code>{stats['user_messages']}</code>\n"
50
+ f"📝 Суммаризировано: <code>{stats['summarized_messages']}</code>\n"
51
+ f"🔢 Токенов: <code>{stats['total_tokens_used']}</code>\n"
52
+ f"⏱ Средняя задержка: <code>{stats['avg_latency_ms']} мс</code>",
53
+ parse_mode="HTML",
54
+ )
55
+
56
+
57
+ @router.callback_query(F.data == "show_settings")
58
+ async def cb_show_settings(callback: types.CallbackQuery) -> None:
59
+ from config import config
60
+ await callback.answer()
61
+ kb = InlineKeyboardMarkup(inline_keyboard=[
62
+ [InlineKeyboardButton(text="🌡 Temperature", callback_data="set_temp")],
63
+ [InlineKeyboardButton(text="📏 Max Tokens", callback_data="set_max_tokens")],
64
+ [InlineKeyboardButton(text="🔙 Назад", callback_data="back_to_menu")],
65
+ ])
66
+ await callback.message.answer(
67
+ "⚙️ <b>Настройки:</b>\n\n"
68
+ f"🌡 Temperature: <code>{config.GLM_TEMPERATURE}</code>\n"
69
+ f"📏 Max Tokens: <code>{config.GLM_MAX_TOKENS}</code>\n"
70
+ f"🔄 Streaming: <code>{'ON' if config.STREAMING_ENABLED else 'OFF'}</code>",
71
+ parse_mode="HTML",
72
+ reply_markup=kb,
73
+ )
74
+
75
+
76
+ @router.callback_query(F.data == "back_to_menu")
77
+ async def cb_back_to_menu(callback: types.CallbackQuery) -> None:
78
+ await callback.answer()
79
+ await callback.message.answer(
80
+ "👋 Главное меню:",
81
+ reply_markup=get_main_menu(),
82
+ )
handlers chat.py ADDED
@@ -0,0 +1,261 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import asyncio
2
+ import logging
3
+ import time
4
+ from aiogram import Router, types, F
5
+ from aiogram.exceptions import TelegramAPIError
6
+ from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton
7
+ from aiogram.enums import ParseMode
8
+
9
+ from config import config
10
+ from database import db
11
+ from services.glm import glm_service, LLMServiceError
12
+ from utils.helpers import send_long_message, edit_or_send_message
13
+
14
+ logger = logging.getLogger(__name__)
15
+ router = Router()
16
+
17
+ # In-memory cache for last user messages (for regeneration)
18
+ _last_user_messages: dict[int, str] = {}
19
+
20
+ SYSTEM_PROMPT = {
21
+ "role": "system",
22
+ "content": (
23
+ "Ты — AI-ассистент на базе GLM-4.5, работающий через NVIDIA NIM API. "
24
+ "Ты помогаешь пользователю как Senior Staff Python Engineer и Acting CTO.\n\n"
25
+ "Правила работы:\n"
26
+ "1. Отвечай на русском языке, если не указано иное.\n"
27
+ "2. Код пиши production-ready, Python 3.11+, async-safe.\n"
28
+ "3. Используй PATCH MODE для изменений >200 строк.\n"
29
+ "4. Не ломай существующую архитектуру.\n"
30
+ "5. Минимальные изменения — локальные патчи.\n"
31
+ "6. Не используй TODO, FIXME, PLACEHOLDER, PASS.\n"
32
+ "7. Проверяй безопасность: SQL Injection, Race Conditions, Async Safety.\n"
33
+ "8. Если не хватает контекста — запроси дополнительный код.\n"
34
+ "9. Не придумывай существующие методы/классы/API.\n"
35
+ "10. Перед кодом: анализ, архитектура, точка интеграции, зависимости."
36
+ ),
37
+ }
38
+
39
+
40
+ def get_chat_buttons() -> InlineKeyboardMarkup:
41
+ """Кнопки под ответом бота."""
42
+ return InlineKeyboardMarkup(inline_keyboard=[
43
+ [
44
+ InlineKeyboardButton(text="🔄 Регенерировать", callback_data="regenerate"),
45
+ InlineKeyboardButton(text="🗑 Очистить", callback_data="clear_memory"),
46
+ ],
47
+ ])
48
+
49
+
50
+ async def _keep_typing(message: types.Message, stop_event: asyncio.Event) -> None:
51
+ """Периодически отправляет typing action, пока stop_event не установлен."""
52
+ while not stop_event.is_set():
53
+ try:
54
+ await message.bot.send_chat_action(message.chat.id, "typing")
55
+ except TelegramAPIError as e:
56
+ logger.debug("Typing action failed: %s", e)
57
+ try:
58
+ await asyncio.wait_for(stop_event.wait(), timeout=config.TYPING_ACTION_INTERVAL)
59
+ except asyncio.TimeoutError:
60
+ continue
61
+
62
+
63
+ async def _maybe_summarize(user_id: int) -> None:
64
+ """Фоновая суммаризация старой истории."""
65
+ try:
66
+ count = await db.count_unsummarized(user_id)
67
+ threshold = config.SUMMARIZE_THRESHOLD
68
+ keep_recent = config.KEEP_RECENT_MESSAGES
69
+
70
+ if count <= threshold:
71
+ return
72
+
73
+ to_summarize = count - keep_recent
74
+ if to_summarize < 5:
75
+ return
76
+
77
+ old_messages = await db.get_oldest_unsummarized(user_id, to_summarize)
78
+ if len(old_messages) < 5:
79
+ return
80
+
81
+ dialog_text = "\n".join([
82
+ f"{m['role']}: {m['content']}" for m in old_messages
83
+ ])
84
+
85
+ summary = await glm_service.summarize(dialog_text)
86
+ if not summary:
87
+ return
88
+
89
+ cutoff_id = old_messages[-1]["id"]
90
+ await db.save_summary(user_id, summary, len(old_messages))
91
+ await db.mark_summarized(user_id, cutoff_id)
92
+
93
+ logger.info(
94
+ "Summarized %d messages for user %s, cutoff_id=%s",
95
+ len(old_messages), user_id, cutoff_id
96
+ )
97
+ except Exception as e:
98
+ logger.error("Summarization failed for user %s: %s", user_id, e, exc_info=True)
99
+
100
+
101
+ async def _build_messages(user_id: int, user_text: str) -> list[dict]:
102
+ """Строит список сообщений для LLM с учётом бюджета токенов."""
103
+ messages: list[dict] = [SYSTEM_PROMPT]
104
+
105
+ # Добавляем summary как контекст
106
+ summary = await db.get_summary(user_id)
107
+ if summary:
108
+ messages.append({
109
+ "role": "system",
110
+ "content": f"[Контекст предыдущих диалогов: {summary}]"
111
+ })
112
+
113
+ # Получаем историю с учётом токенового бюджета
114
+ history = await db.get_messages_with_token_budget(
115
+ user_id, max_tokens=config.MAX_CONTEXT_TOKENS
116
+ )
117
+
118
+ # Ограничиваем количество сообщений
119
+ if len(history) > config.MAX_HISTORY:
120
+ history = history[-config.MAX_HISTORY:]
121
+
122
+ for h in history:
123
+ messages.append({"role": h["role"], "content": h["content"]})
124
+
125
+ # Добавляем текущее сообщение (если его ещё нет в истории)
126
+ if not history or history[-1]["content"] != user_text:
127
+ messages.append({"role": "user", "content": user_text})
128
+
129
+ return messages
130
+
131
+
132
+ @router.message(F.text)
133
+ async def handle_chat(message: types.Message) -> None:
134
+ """Основной обработчик текстовых сообщений."""
135
+ logger.info("→ handle_chat START, user_id=%s, text_len=%d", message.from_user.id, len(message.text or ""))
136
+
137
+ if not message.text:
138
+ await message.answer("❌ Поддерживаются только текстовые сообщения.")
139
+ return
140
+
141
+ user_id = message.from_user.id
142
+ user_text = message.text
143
+ _last_user_messages[user_id] = user_text
144
+
145
+ # Upsert user
146
+ await db.upsert_user(
147
+ user_id,
148
+ message.from_user.username,
149
+ message.from_user.first_name,
150
+ message.from_user.last_name,
151
+ )
152
+
153
+ # Save user message
154
+ tokens_estimate = glm_service.estimate_tokens(user_text)
155
+ await db.save_message(user_id, "user", user_text, tokens_estimate)
156
+
157
+ # Build messages for LLM
158
+ messages = await _build_messages(user_id, user_text)
159
+
160
+ # Typing action
161
+ stop_typing = asyncio.Event()
162
+ typing_task = asyncio.create_task(_keep_typing(message, stop_typing))
163
+
164
+ start_time = time.perf_counter()
165
+ bot_message = None
166
+
167
+ try:
168
+ logger.info("→ Calling GLM, messages_count=%d", len(messages))
169
+
170
+ if config.STREAMING_ENABLED:
171
+ # Streaming mode — обновляем сообщение по мере поступления чанков
172
+ response_text = ""
173
+ chunk_buffer = ""
174
+ last_update = time.perf_counter()
175
+ message_sent = False
176
+
177
+ async for chunk in glm_service.chat_stream(messages, user_id=user_id):
178
+ response_text += chunk
179
+ chunk_buffer += chunk
180
+
181
+ # Обновляем сообщение не чаще чем раз в N секунд
182
+ now = time.perf_counter()
183
+ if now - last_update >= config.STREAMING_UPDATE_INTERVAL and len(response_text) > 10:
184
+ if not message_sent:
185
+ bot_message = await message.answer(
186
+ response_text + "▌",
187
+ parse_mode=ParseMode.MARKDOWN,
188
+ )
189
+ message_sent = True
190
+ else:
191
+ try:
192
+ await bot_message.edit_text(
193
+ response_text + "▌",
194
+ parse_mode=ParseMode.MARKDOWN,
195
+ )
196
+ except TelegramAPIError:
197
+ pass # Message not modified или другая ошибка
198
+ last_update = now
199
+ chunk_buffer = ""
200
+
201
+ # Финальное обновление
202
+ stop_typing.set()
203
+ await typing_task
204
+
205
+ if message_sent and bot_message:
206
+ try:
207
+ await bot_message.edit_text(
208
+ response_text,
209
+ parse_mode=ParseMode.MARKDOWN,
210
+ reply_markup=get_chat_buttons(),
211
+ )
212
+ except TelegramAPIError:
213
+ # Если edit не удался, отправляем новое
214
+ await send_long_message(
215
+ message, response_text,
216
+ reply_markup=get_chat_buttons(),
217
+ )
218
+ else:
219
+ await send_long_message(
220
+ message, response_text,
221
+ reply_markup=get_chat_buttons(),
222
+ )
223
+
224
+ else:
225
+ # Non-streaming mode
226
+ response = await glm_service.chat(messages, user_id=user_id)
227
+ stop_typing.set()
228
+ await typing_task
229
+ await send_long_message(message, response, reply_markup=get_chat_buttons())
230
+ response_text = response
231
+
232
+ duration_ms = (time.perf_counter() - start_time) * 1000
233
+ logger.info("← Response sent, duration=%.1fms, length=%d", duration_ms, len(response_text))
234
+
235
+ # Save assistant message
236
+ assistant_tokens = glm_service.estimate_tokens(response_text)
237
+ await db.save_message(user_id, "assistant", response_text, assistant_tokens)
238
+
239
+ # Trigger background summarization
240
+ asyncio.create_task(_maybe_summarize(user_id))
241
+
242
+ except LLMServiceError as e:
243
+ stop_typing.set()
244
+ await typing_task
245
+ logger.error("LLM Service Error: %s", e)
246
+ await message.answer(
247
+ "❌ <b>Ошибка модели:</b>\n"
248
+ f"<code>{str(e)[:200]}</code>\n\n"
249
+ "Попробуйте повторить запрос позже.",
250
+ parse_mode="HTML",
251
+ )
252
+
253
+ except Exception as e:
254
+ stop_typing.set()
255
+ await typing_task
256
+ logger.error("❌ Unexpected error in chat handler: %s", e, exc_info=True)
257
+ await message.answer(
258
+ "❌ <b>Неожиданная ошибка.</b>\n"
259
+ "Разработчик уведомлён. Попробуйте позже.",
260
+ parse_mode="HTML",
261
+ )
handlers commands.py ADDED
@@ -0,0 +1,149 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import asyncio
2
+ import logging
3
+ from aiogram import Router, types, F
4
+ from aiogram.filters import Command
5
+ from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton
6
+ from database import db
7
+ from config import config
8
+
9
+ logger = logging.getLogger(__name__)
10
+ router = Router()
11
+
12
+
13
+ def get_main_menu() -> InlineKeyboardMarkup:
14
+ """Главное меню с быстрыми действиями."""
15
+ return InlineKeyboardMarkup(inline_keyboard=[
16
+ [
17
+ InlineKeyboardButton(text="🔄 Регенерировать", callback_data="regenerate"),
18
+ InlineKeyboardButton(text="🗑 Очистить память", callback_data="clear_memory"),
19
+ ],
20
+ [
21
+ InlineKeyboardButton(text="📊 Статистика", callback_data="show_stats"),
22
+ InlineKeyboardButton(text="⚙️ Настройки", callback_data="show_settings"),
23
+ ],
24
+ ])
25
+
26
+
27
+ @router.message(Command("start"))
28
+ async def cmd_start(message: types.Message) -> None:
29
+ await message.answer(
30
+ "👋 <b>Привет!</b> Я AI-ассистент на базе GLM-4.5 через NVIDIA NIM.\n\n"
31
+ "🧠 <b>Возможности:</b>\n"
32
+ "• Долговременная память (Neon PostgreSQL)\n"
33
+ "• Автоматическая суммаризация диалогов\n"
34
+ "• Fallback-модель при сбоях\n"
35
+ "• Streaming-ответы\n"
36
+ "• Rate limiting и мониторинг\n\n"
37
+ "Отправь любое сообщение — я отвечу.\n"
38
+ "Используй /help для списка команд.",
39
+ parse_mode="HTML",
40
+ reply_markup=get_main_menu(),
41
+ )
42
+
43
+
44
+ @router.message(Command("help"))
45
+ async def cmd_help(message: types.Message) -> None:
46
+ await message.answer(
47
+ "📋 <b>Команды:</b>\n\n"
48
+ "/start — Начать работу\n"
49
+ "/help — Помощь\n"
50
+ "/clear — Очистить историю сообщений\n"
51
+ "/stats — Статистика использования\n"
52
+ "/memory — Информация о памяти\n"
53
+ "/settings — Настройки бота\n"
54
+ "/ping — Проверка задержки\n\n"
55
+ "<b>Быстрые кнопки под ответом:</b>\n"
56
+ "🔄 Регенерировать — повторить генерацию\n"
57
+ "🗑 Очистить память — сбросить контекст\n"
58
+ "📊 Статистика — метрики и usage",
59
+ parse_mode="HTML",
60
+ )
61
+
62
+
63
+ @router.message(Command("clear"))
64
+ async def cmd_clear(message: types.Message) -> None:
65
+ count = await db.clear_history(message.from_user.id)
66
+ await message.answer(
67
+ f"🗑 <b>История очищена.</b>\n\n"
68
+ f"Удалено сообщений: <code>{count}</code>.\n"
69
+ f"Контекст и суммаризация сброшены.",
70
+ parse_mode="HTML",
71
+ )
72
+
73
+
74
+ @router.message(Command("stats"))
75
+ async def cmd_stats(message: types.Message) -> None:
76
+ stats = await db.get_stats(message.from_user.id)
77
+ await message.answer(
78
+ f"📊 <b>Статистика:</b>\n\n"
79
+ f"👥 Всего пользователей: <code>{stats['total_users']}</code>\n"
80
+ f"💬 Ваших сообщений: <code>{stats['user_messages']}</code>\n"
81
+ f"📨 Всего сообщений: <code>{stats['total_messages']}</code>\n"
82
+ f"📝 Суммаризировано: <code>{stats['summarized_messages']}</code>\n"
83
+ f"🔢 Токенов использовано: <code>{stats['total_tokens_used']}</code>\n"
84
+ f"⏱ Средняя задержка: <code>{stats['avg_latency_ms']} мс</code>",
85
+ parse_mode="HTML",
86
+ )
87
+
88
+
89
+ @router.message(Command("memory"))
90
+ async def cmd_memory(message: types.Message) -> None:
91
+ user_id = message.from_user.id
92
+ summary = await db.get_summary(user_id)
93
+ count = await db.count_unsummarized(user_id)
94
+ history = await db.get_messages(user_id, limit=5)
95
+
96
+ text = "🧠 <b>Память пользователя:</b>\n\n"
97
+
98
+ if summary:
99
+ text += f"📋 <b>Суммаризация:</b>\n<code>{summary[:500]}</code>"
100
+ if len(summary) > 500:
101
+ text += "..."
102
+ text += "\n\n"
103
+ else:
104
+ text += "📋 <b>Суммаризация:</b> <i>пока нет</i>\n\n"
105
+
106
+ text += f"💬 <b>Несуммаризированных сообщений:</b> <code>{count}</code>\n"
107
+ text += f"📚 <b>Лимит истории:</b> <code>{config.MAX_HISTORY}</code>\n"
108
+ text += f"🔢 <b>Лимит токенов:</b> <code>{config.MAX_CONTEXT_TOKENS}</code>\n\n"
109
+
110
+ if history:
111
+ text += "<b>Последние сообщения:</b>\n"
112
+ for h in history[-3:]:
113
+ preview = h['content'][:60].replace("<", "&lt;").replace(">", "&gt;")
114
+ text += f"• <i>{h['role']}</i>: {preview}...\n"
115
+
116
+ await message.answer(text, parse_mode="HTML")
117
+
118
+
119
+ @router.message(Command("settings"))
120
+ async def cmd_settings(message: types.Message) -> None:
121
+ kb = InlineKeyboardMarkup(inline_keyboard=[
122
+ [InlineKeyboardButton(text="🌡 Temperature", callback_data="set_temp")],
123
+ [InlineKeyboardButton(text="📏 Max Tokens", callback_data="set_max_tokens")],
124
+ [InlineKeyboardButton(text="🔙 Назад", callback_data="back_to_menu")],
125
+ ])
126
+ await message.answer(
127
+ "⚙️ <b>Настройки:</b>\n\n"
128
+ f"🌡 <b>Temperature:</b> <code>{config.GLM_TEMPERATURE}</code>\n"
129
+ f"📏 <b>Max Tokens:</b> <code>{config.GLM_MAX_TOKENS}</code>\n"
130
+ f"🔄 <b>Streaming:</b> <code>{'ON' if config.STREAMING_ENABLED else 'OFF'}</code>\n"
131
+ f"🛡 <b>Rate Limit:</b> <code>{config.RATE_LIMIT_REQUESTS_PER_MINUTE}/мин</code>\n"
132
+ f"💾 <b>Cache:</b> <code>{'ON' if config.CACHE_ENABLED else 'OFF'}</code>",
133
+ parse_mode="HTML",
134
+ reply_markup=kb,
135
+ )
136
+
137
+
138
+ @router.message(Command("ping"))
139
+ async def cmd_ping(message: types.Message) -> None:
140
+ start = asyncio.get_event_loop().time()
141
+ msg = await message.answer("🏓 <b>Pong!</b>", parse_mode="HTML")
142
+ end = asyncio.get_event_loop().time()
143
+ latency = (end - start) * 1000
144
+ await msg.edit_text(
145
+ f"🏓 <b>Pong!</b>\n\n"
146
+ f"⏱ Задержка Telegram: <code>{latency:.1f} мс</code>\n"
147
+ f"🎯 Бот активен и готов к работе.",
148
+ parse_mode="HTML",
149
+ )
middlewares rate_limit.py ADDED
@@ -0,0 +1,43 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import logging
2
+ from typing import Any, Awaitable, Callable, Dict
3
+ from aiogram import BaseMiddleware
4
+ from aiogram.types import Message, TelegramObject
5
+ from database import db
6
+ from config import config
7
+
8
+ logger = logging.getLogger(__name__)
9
+
10
+
11
+ class RateLimitMiddleware(BaseMiddleware):
12
+ """
13
+ Middleware для ограничения частоты запросов.
14
+ Хранит состояние в PostgreSQL (stateless — работает на нескольких инстансах).
15
+ """
16
+
17
+ async def __call__(
18
+ self,
19
+ handler: Callable[[TelegramObject, Dict[str, Any]], Awaitable[Any]],
20
+ event: TelegramObject,
21
+ data: Dict[str, Any],
22
+ ) -> Any:
23
+ if not isinstance(event, Message):
24
+ return await handler(event, data)
25
+
26
+ if not config.RATE_LIMIT_ENABLED:
27
+ return await handler(event, data)
28
+
29
+ user_id = event.from_user.id
30
+ allowed, remaining, reset_in = await db.check_rate_limit(user_id)
31
+
32
+ if not allowed:
33
+ logger.warning("Rate limit exceeded for user %s", user_id)
34
+ minutes = int(reset_in // 60) + 1
35
+ await event.answer(
36
+ f"⏳ Слишком много запросов. Попробуйте снова через {minutes} мин.",
37
+ show_alert=False,
38
+ )
39
+ return None
40
+
41
+ # Добавляем информацию о rate limit в data для логирования
42
+ data["rate_limit_remaining"] = remaining
43
+ return await handler(event, data)
services glm.py ADDED
@@ -0,0 +1,344 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import asyncio
2
+ import logging
3
+ import time
4
+ from typing import List, Dict, Any, Optional, AsyncGenerator
5
+
6
+ import httpx
7
+ import openai
8
+ from openai import AsyncOpenAI
9
+
10
+ from config import config
11
+ from database import db
12
+
13
+ logger = logging.getLogger(__name__)
14
+
15
+
16
+ class LLMServiceError(Exception):
17
+ """Base exception for LLM service errors."""
18
+ pass
19
+
20
+
21
+ class LLMTimeoutError(LLMServiceError):
22
+ """Timeout error with fallback exhausted."""
23
+ pass
24
+
25
+
26
+ class LLMRateLimitError(LLMServiceError):
27
+ """Rate limit from provider."""
28
+ pass
29
+
30
+
31
+ class GLMService:
32
+ """Production-ready GLM service with retry, fallback, streaming, metrics."""
33
+
34
+ def __init__(self) -> None:
35
+ # Granular timeouts for NVIDIA NIM
36
+ self.timeout = httpx.Timeout(
37
+ connect=config.TIMEOUT_CONNECT,
38
+ read=config.TIMEOUT_READ,
39
+ write=config.TIMEOUT_WRITE,
40
+ pool=config.TIMEOUT_POOL,
41
+ )
42
+
43
+ self.primary_client = AsyncOpenAI(
44
+ base_url=config.NVIDIA_BASE_URL,
45
+ api_key=config.NVIDIA_API_KEY,
46
+ timeout=self.timeout,
47
+ max_retries=0, # управляем retry самостоятельно
48
+ )
49
+
50
+ self.fallback_client = AsyncOpenAI(
51
+ base_url=config.NVIDIA_BASE_URL,
52
+ api_key=config.NVIDIA_API_KEY,
53
+ timeout=self.timeout,
54
+ max_retries=0,
55
+ ) if config.FALLBACK_ENABLED else None
56
+
57
+ self.primary_model = config.PRIMARY_MODEL
58
+ self.fallback_model = config.FALLBACK_MODEL
59
+
60
+ # ═══════════════════════════════════════════════════════════════
61
+ # BLOCK: Exponential Retry with Jitter
62
+ # ═══════════════════════════════════════════════════════════════
63
+ def _calculate_delay(self, attempt: int) -> float:
64
+ """Экспоненциальная задержка с jitter."""
65
+ import random
66
+ delay = min(
67
+ config.RETRY_BASE_DELAY * (config.RETRY_EXPONENTIAL_BASE ** attempt),
68
+ config.RETRY_MAX_DELAY,
69
+ )
70
+ jitter = random.uniform(0, delay * 0.3)
71
+ return delay + jitter
72
+
73
+ def _is_retryable_error(self, error: Exception) -> bool:
74
+ """Определяет, стоит ли retry-ить ошибку."""
75
+ if isinstance(error, openai.APITimeoutError):
76
+ return True
77
+ if isinstance(error, openai.APIConnectionError):
78
+ return True
79
+ if isinstance(error, openai.RateLimitError):
80
+ return True
81
+ if isinstance(error, openai.InternalServerError):
82
+ return True
83
+ if isinstance(error, openai.APIStatusError):
84
+ if hasattr(error, 'status_code') and error.status_code in (429, 502, 503, 504):
85
+ return True
86
+ return False
87
+
88
+ # ═══════════════════════════════════════════════════════════════
89
+ # BLOCK: Core Chat with Retry and Fallback
90
+ # ═══════════════════════════════════════════════════════════════
91
+ async def chat(
92
+ self,
93
+ messages: List[Dict[str, str]],
94
+ user_id: Optional[int] = None,
95
+ temperature: Optional[float] = None,
96
+ max_tokens: Optional[int] = None,
97
+ stream: bool = False,
98
+ ) -> str:
99
+ """
100
+ Отправляет запрос к LLM с retry и fallback.
101
+ Возвращает полный текст ответа.
102
+ """
103
+ start_time = time.perf_counter()
104
+ model = self.primary_model
105
+ client = self.primary_client
106
+ used_fallback = False
107
+ last_error = None
108
+
109
+ params = self._build_params(
110
+ messages=messages,
111
+ model=model,
112
+ temperature=temperature,
113
+ max_tokens=max_tokens,
114
+ stream=stream,
115
+ )
116
+
117
+ # Retry loop for primary model
118
+ for attempt in range(config.MAX_RETRIES):
119
+ try:
120
+ logger.info(
121
+ "→ LLM request: model=%s, attempt=%d/%d, messages=%d, stream=%s",
122
+ model, attempt + 1, config.MAX_RETRIES, len(messages), stream
123
+ )
124
+
125
+ if stream and config.STREAMING_ENABLED:
126
+ response_text = await self._stream_chat(client, params, user_id)
127
+ else:
128
+ response = await client.chat.completions.create(**params)
129
+ response_text = response.choices[0].message.content or ""
130
+
131
+ duration_ms = (time.perf_counter() - start_time) * 1000
132
+ logger.info(
133
+ "← LLM response: model=%s, duration=%.1fms, length=%d, fallback=%s",
134
+ model, duration_ms, len(response_text), used_fallback
135
+ )
136
+
137
+ # Save metrics
138
+ if config.METRICS_ENABLED:
139
+ await db.save_metric(
140
+ user_id=user_id,
141
+ model=model,
142
+ duration_ms=duration_ms,
143
+ success=True,
144
+ )
145
+
146
+ return response_text
147
+
148
+ except Exception as e:
149
+ last_error = e
150
+ duration_ms = (time.perf_counter() - start_time) * 1000
151
+ error_type = type(e).__name__
152
+
153
+ if not self._is_retryable_error(e):
154
+ logger.error("Non-retryable error: %s", e)
155
+ if config.METRICS_ENABLED:
156
+ await db.save_metric(
157
+ user_id=user_id,
158
+ model=model,
159
+ duration_ms=duration_ms,
160
+ success=False,
161
+ error_type=error_type,
162
+ )
163
+ raise LLMServiceError(f"Non-retryable error: {e}") from e
164
+
165
+ if attempt < config.MAX_RETRIES - 1:
166
+ delay = self._calculate_delay(attempt)
167
+ logger.warning(
168
+ "⚠️ Retry %d/%d for model=%s after %.1fs: %s",
169
+ attempt + 1, config.MAX_RETRIES, model, delay, e
170
+ )
171
+ await asyncio.sleep(delay)
172
+ else:
173
+ logger.error("Primary model exhausted all retries: %s", e)
174
+
175
+ # Fallback model
176
+ if config.FALLBACK_ENABLED and self.fallback_client:
177
+ logger.warning("🔄 Switching to fallback model: %s", self.fallback_model)
178
+ model = self.fallback_model
179
+ client = self.fallback_client
180
+ used_fallback = True
181
+
182
+ try:
183
+ params["model"] = model
184
+ if stream and config.STREAMING_ENABLED:
185
+ response_text = await self._stream_chat(client, params, user_id)
186
+ else:
187
+ response = await client.chat.completions.create(**params)
188
+ response_text = response.choices[0].message.content or ""
189
+
190
+ duration_ms = (time.perf_counter() - start_time) * 1000
191
+ logger.info(
192
+ "← Fallback response: model=%s, duration=%.1fms, length=%d",
193
+ model, duration_ms, len(response_text)
194
+ )
195
+
196
+ if config.METRICS_ENABLED:
197
+ await db.save_metric(
198
+ user_id=user_id,
199
+ model=f"{model} (fallback)",
200
+ duration_ms=duration_ms,
201
+ success=True,
202
+ )
203
+
204
+ return response_text
205
+
206
+ except Exception as e:
207
+ duration_ms = (time.perf_counter() - start_time) * 1000
208
+ logger.error("Fallback model also failed: %s", e)
209
+ if config.METRICS_ENABLED:
210
+ await db.save_metric(
211
+ user_id=user_id,
212
+ model=f"{model} (fallback)",
213
+ duration_ms=duration_ms,
214
+ success=False,
215
+ error_type=type(e).__name__,
216
+ )
217
+ raise LLMTimeoutError(
218
+ f"Both primary and fallback models failed. Last error: {last_error}"
219
+ ) from e
220
+
221
+ raise LLMTimeoutError(f"All retries exhausted. Last error: {last_error}")
222
+
223
+ # ═══════════════════════════════════════════════════════════════
224
+ # BLOCK: Streaming
225
+ # ═══════════════════════════════════════════════════════════════
226
+ async def _stream_chat(
227
+ self,
228
+ client: AsyncOpenAI,
229
+ params: Dict[str, Any],
230
+ user_id: Optional[int] = None,
231
+ ) -> str:
232
+ """Обрабатывает streaming-ответ и собирает полный текст."""
233
+ params["stream"] = True
234
+ params["stream_options"] = {"include_usage": True}
235
+
236
+ full_text = ""
237
+ usage = None
238
+
239
+ try:
240
+ stream = await client.chat.completions.create(**params)
241
+ async for chunk in stream:
242
+ if chunk.choices and chunk.choices[0].delta.content:
243
+ full_text += chunk.choices[0].delta.content
244
+ if chunk.usage:
245
+ usage = chunk.usage
246
+ except Exception as e:
247
+ logger.error("Streaming error: %s", e)
248
+ raise
249
+
250
+ if usage:
251
+ logger.info(
252
+ "Streaming complete: tokens_in=%d, tokens_out=%d",
253
+ usage.prompt_tokens or 0, usage.completion_tokens or 0
254
+ )
255
+
256
+ return full_text
257
+
258
+ async def chat_stream(
259
+ self,
260
+ messages: List[Dict[str, str]],
261
+ user_id: Optional[int] = None,
262
+ ) -> AsyncGenerator[str, None]:
263
+ """Yields text chunks for real-time Telegram updates."""
264
+ params = self._build_params(
265
+ messages=messages,
266
+ model=self.primary_model,
267
+ stream=True,
268
+ )
269
+ params["stream"] = True
270
+ params["stream_options"] = {"include_usage": True}
271
+
272
+ try:
273
+ stream = await self.primary_client.chat.completions.create(**params)
274
+ async for chunk in stream:
275
+ if chunk.choices and chunk.choices[0].delta.content:
276
+ yield chunk.choices[0].delta.content
277
+ except Exception as e:
278
+ logger.error("Streaming generation failed: %s", e)
279
+ raise
280
+
281
+ # ═══════════════════════════════════════════════════════════════
282
+ # BLOCK: Summarization
283
+ # ═══════════════════════════════════════════════════════════════
284
+ async def summarize(self, dialog_text: str) -> str:
285
+ """Суммаризация диалога через LLM."""
286
+ start_time = time.perf_counter()
287
+ messages = [
288
+ {
289
+ "role": "system",
290
+ "content": (
291
+ "Суммаризируй следующий диалог между пользователем и ассистентом. "
292
+ "Сохрани ключевые факты, предпочтения пользователя, важные детали и контекст. "
293
+ "Будь краток, максимум 4096 токенов. Используй русский язык."
294
+ )
295
+ },
296
+ {"role": "user", "content": dialog_text}
297
+ ]
298
+
299
+ try:
300
+ response = await self.primary_client.chat.completions.create(
301
+ model=self.primary_model,
302
+ messages=messages,
303
+ temperature=0.1,
304
+ max_tokens=config.SUMMARY_MAX_TOKENS,
305
+ timeout=httpx.Timeout(connect=10.0, read=60.0, write=10.0),
306
+ )
307
+ content = response.choices[0].message.content or ""
308
+ duration_ms = (time.perf_counter() - start_time) * 1000
309
+ logger.info("Summary generated in %.1fms, length=%d", duration_ms, len(content))
310
+ return content
311
+ except Exception as e:
312
+ logger.error("Summary generation failed: %s", e)
313
+ return ""
314
+
315
+ # ═══════════════════════════════════════════════════════════════
316
+ # BLOCK: Helpers
317
+ # ═══════════════════════════════════════════════════════════════
318
+ def _build_params(
319
+ self,
320
+ messages: List[Dict[str, str]],
321
+ model: str,
322
+ temperature: Optional[float] = None,
323
+ max_tokens: Optional[int] = None,
324
+ stream: bool = False,
325
+ ) -> Dict[str, Any]:
326
+ """Строит параметры запроса к API."""
327
+ return {
328
+ "model": model,
329
+ "messages": messages,
330
+ "temperature": temperature if temperature is not None else config.GLM_TEMPERATURE,
331
+ "top_p": config.GLM_TOP_P,
332
+ "frequency_penalty": config.GLM_FREQUENCY_PENALTY,
333
+ "presence_penalty": config.GLM_PRESENCE_PENALTY,
334
+ "max_tokens": max_tokens if max_tokens is not None else config.GLM_MAX_TOKENS,
335
+ "stream": stream,
336
+ }
337
+
338
+ def estimate_tokens(self, text: str) -> int:
339
+ """Грубая оценка количества токенов (1 token ≈ 0.75 английских слов или 0.5 русских)."""
340
+ # Простая эвристика: ~4 символа на токен для смешанного текста
341
+ return max(1, len(text) // 4)
342
+
343
+
344
+ glm_service = GLMService()
utils helpers.py ADDED
@@ -0,0 +1,114 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import logging
2
+ import re
3
+ from aiogram import types
4
+ from aiogram.enums import ParseMode
5
+ from config import config
6
+
7
+ logger = logging.getLogger(__name__)
8
+ MAX_LENGTH = config.MAX_MESSAGE_LENGTH
9
+
10
+
11
+ def escape_markdown_v2(text: str) -> str:
12
+ """Экранирует спецсимволы для MarkdownV2 (Telegram)."""
13
+ # Символы, которые нужно экранировать в MarkdownV2
14
+ escape_chars = r"_\*\[\]\(\)~`>#+\-=|{}\.!"
15
+ return re.sub(f"([{re.escape(escape_chars)}])", r"\\1", text)
16
+
17
+
18
+ def format_code_blocks(text: str) -> str:
19
+ """Гарантирует корректное закрытие code blocks."""
20
+ # Подсчёт открывающих и закрывающих ```
21
+ triple_backticks = text.count("```")
22
+ if triple_backticks % 2 != 0:
23
+ text += "
24
+ ```"
25
+ return text
26
+
27
+
28
+ def split_message_smart(text: str, max_length: int = MAX_LENGTH) -> list[str]:
29
+ """
30
+ Разбивает текст на части, стараясь не ломать code blocks и параграфы.
31
+ """
32
+ if len(text) <= max_length:
33
+ return [text]
34
+
35
+ chunks = []
36
+ while text:
37
+ if len(text) <= max_length:
38
+ chunks.append(text)
39
+ break
40
+
41
+ # Ищем ближайший перенос строки перед max_length
42
+ split_at = text.rfind("\n", 0, max_length)
43
+ if split_at == -1:
44
+ split_at = text.rfind(" ", 0, max_length)
45
+ if split_at == -1:
46
+ split_at = max_length
47
+
48
+ chunk = text[:split_at].strip()
49
+ if chunk:
50
+ chunks.append(chunk)
51
+ text = text[split_at:].strip()
52
+
53
+ return chunks
54
+
55
+
56
+ async def send_long_message(
57
+ message: types.Message,
58
+ text: str,
59
+ parse_mode: ParseMode = ParseMode.MARKDOWN,
60
+ reply_markup=None,
61
+ ) -> list[types.Message]:
62
+ """
63
+ Отправляет длинные сообщения частями с корректным Markdown.
64
+ Возвращает список отправленных сообщений.
65
+ """
66
+ text = format_code_blocks(text)
67
+ chunks = split_message_smart(text)
68
+ sent_messages = []
69
+
70
+ for i, chunk in enumerate(chunks):
71
+ try:
72
+ # Для последнего чанка добавляем reply_markup (кнопки)
73
+ markup = reply_markup if i == len(chunks) - 1 else None
74
+ sent = await message.answer(
75
+ chunk,
76
+ parse_mode=parse_mode,
77
+ reply_markup=markup,
78
+ )
79
+ sent_messages.append(sent)
80
+ except Exception as e:
81
+ logger.warning("Failed to send chunk with Markdown: %s. Sending as plain text.", e)
82
+ try:
83
+ sent = await message.answer(
84
+ chunk,
85
+ parse_mode=None,
86
+ reply_markup=reply_markup if i == len(chunks) - 1 else None,
87
+ )
88
+ sent_messages.append(sent)
89
+ except Exception as e2:
90
+ logger.error("Failed to send chunk: %s", e2)
91
+
92
+ return sent_messages
93
+
94
+
95
+ async def edit_or_send_message(
96
+ message: types.Message,
97
+ text: str,
98
+ last_bot_message: types.Message = None,
99
+ parse_mode: ParseMode = ParseMode.MARKDOWN,
100
+ ) -> types.Message:
101
+ """
102
+ Редактирует последнее сообщение бота или отправляет новое.
103
+ Используется для streaming-ответов.
104
+ """
105
+ text = format_code_blocks(text)
106
+
107
+ if last_bot_message and len(text) <= MAX_LENGTH:
108
+ try:
109
+ await last_bot_message.edit_text(text, parse_mode=parse_mode)
110
+ return last_bot_message
111
+ except Exception as e:
112
+ logger.debug("Edit failed, sending new message: %s", e)
113
+
114
+ return await message.answer(text, parse_mode=parse_mode)