main
FastAPI-приложение навык LostAndFound для Яндекс Алисы.
Основная точка входа: обрабатывает POST-запросы от Яндекс Диалогов, парсит команды пользователя и выполняет CRUD-операции с местоположения вещей и возвращает ответы в формате JSON. Также предоставляет health-check и фоновые ассистенты (очистка Langfuse) через APScheduler.
1""" 2FastAPI-приложение навык LostAndFound для Яндекс Алисы. 3 4Основная точка входа: обрабатывает POST-запросы от Яндекс Диалогов, 5парсит команды пользователя и выполняет CRUD-операции с местоположения вещей 6и возвращает ответы в формате JSON. Также предоставляет health-check 7и фоновые ассистенты (очистка Langfuse) через APScheduler. 8""" 9 10import asyncio 11import logging 12import os 13import sys 14import uuid 15from contextlib import asynccontextmanager 16from datetime import datetime 17from typing import Optional 18 19from apscheduler.schedulers.asyncio import AsyncIOScheduler 20from apscheduler.triggers.cron import CronTrigger 21from e_lib import Config 22from e_lib.assistants import LangfuseCleanerAssistant, LLMWarmupAssistant 23from fastapi import APIRouter, FastAPI, HTTPException 24from fastapi.staticfiles import StaticFiles 25from langfuse import Langfuse, observe 26from loguru import logger 27 28from command_parser import CommandParser 29from models import AliceRequest, AppConfig 30from thingman import DB_FILE, ThingManager 31 32APP_VERSION = os.getenv("PRODUCTS_VERSION", "dev") 33 34command_parser = CommandParser() 35thing_manager = ThingManager() 36 37_langfuse_instance: Optional[Langfuse] = None 38 39llm_warmup_assistant: Optional[LLMWarmupAssistant] = None 40_llm_warmup_tasks: set = set() 41 42 43def _trigger_llm_warmup() -> None: 44 """ 45 Запускает фоновый прогрев LLM (best-effort). 46 47 Вызывается при старте сессии навыка: пока пользователь формулирует 48 вопрос, ассистент подгружает модель лёгким запросом. Повторные вызовы 49 при активном прогоне пропускаются самим ассистентом. Ошибки не роняют 50 запрос — прогрев строго фоновый. 51 """ 52 if llm_warmup_assistant is None: 53 return 54 task = asyncio.create_task(llm_warmup_assistant.run()) 55 _llm_warmup_tasks.add(task) 56 task.add_done_callback(_llm_warmup_tasks.discard) 57 58 59def get_langfuse() -> Langfuse: 60 """ 61 Возвращает экземпляр Langfuse для observability (синглтон). 62 63 При первом вызове создаёт клиент с учётом переменных окружения 64 LANGFUSE_PUBLIC_KEY, LANGFUSE_SECRET_KEY и LANGFUSE_HOST. 65 66 Returns: 67 Langfuse: Экземпляр Langfuse-клиента 68 """ 69 global _langfuse_instance 70 if _langfuse_instance is None: 71 _langfuse_instance = Langfuse( 72 public_key=os.getenv("LANGFUSE_PUBLIC_KEY"), 73 secret_key=os.getenv("LANGFUSE_SECRET_KEY"), 74 host=os.getenv("LANGFUSE_HOST"), 75 ) 76 return _langfuse_instance 77 78 79config = Config(os.path.join(os.path.dirname(__file__), "config.yaml"), model=AppConfig) 80 81 82class InterceptHandler(logging.Handler): 83 """ 84 Перехватчик логов из стандартного logging в loguru. 85 86 Перенаправляет все сообщения из логгеров uvicorn и fastapi 87 в форматированный вывод loguru. 88 """ 89 90 def emit(self, record): 91 """ 92 Обрабатывает и перенаправляет запись лога в loguru. 93 94 Args: 95 record: Запись лога из стандартного модуля logging 96 """ 97 try: 98 level = logger.level(record.levelname).name 99 except ValueError: 100 level = record.levelno 101 102 frame, depth = logging.currentframe(), 2 103 while frame and frame.f_code.co_filename == logging.__file__: 104 frame = frame.f_back 105 depth += 1 106 107 logger.opt(depth=depth, exception=record.exc_info).log( 108 level, record.getMessage() 109 ) 110 111 112def setup_logging(): 113 """ 114 Настраивает логирование приложения. 115 116 Перенаправляет логи uvicorn и fastapi в loguru с цветным 117 форматированным выводом в stdout. Уровень DEBUG для всех 118 сообщений приложения. 119 """ 120 logger.remove() 121 122 logger.add( 123 sys.stdout, 124 format="<green>{time:YYYY-MM-DD HH:mm:ss}</green> | <level>{level: <8}</level> | <cyan>{name}</cyan>:<cyan>{function}</cyan>:<cyan>{line}</cyan> | <level>{message}</level>", 125 level="DEBUG", 126 colorize=True, 127 ) 128 129 intercept_handler = InterceptHandler() 130 131 root_logger = logging.getLogger() 132 root_logger.handlers = [intercept_handler] 133 root_logger.setLevel(logging.INFO) 134 135 for logger_name in ["uvicorn", "uvicorn.access", "uvicorn.error", "fastapi"]: 136 logging_logger = logging.getLogger(logger_name) 137 logging_logger.handlers = [intercept_handler] 138 logging_logger.setLevel(logging.INFO) 139 logging_logger.propagate = False 140 141 142@observe(name="Build response") 143def build_response(request: AliceRequest, text: str, tts: str | None = None) -> dict: 144 """ 145 Формирует ответ навыка Яндекс Алисы. 146 147 Args: 148 request: Входящий запрос Алиса (для копирования session/version) 149 text: Текстовый ответ навыка 150 tts: Озвучивемый текст (если отличается от text) 151 152 Returns: 153 dict: Ответ в формате Яндекс Диалогов 154 """ 155 response: dict = { 156 "text": text, 157 "end_session": False, 158 } 159 if tts is not None: 160 response["tts"] = tts 161 return { 162 "version": request.version, 163 "session": { 164 "message_id": request.session.message_id, 165 "session_id": request.session.session_id, 166 "user_id": request.session.user_id, 167 "skill_id": request.session.skill_id, 168 }, 169 "response": response, 170 } 171 172 173@observe(name="Build error response") 174def build_error_response( 175 text: str = "Произошла ошибка при обработке запроса. Попробуйте еще раз.", 176) -> dict: 177 """ 178 Формирует ответ с сообщением об ошибке. 179 180 Args: 181 text: Текст ошибки (по умолчанию стандартное сообщение) 182 183 Returns: 184 dict: Ответ-заглушка с текстом ошибки 185 """ 186 return { 187 "version": "1.0", 188 "session": { 189 "message_id": 0, 190 "session_id": "", 191 "user_id": "", 192 "skill_id": "", 193 }, 194 "response": { 195 "text": text, 196 "end_session": False, 197 }, 198 } 199 200 201web_router = APIRouter() 202 203 204@web_router.get("/api/locations") 205async def web_list_locations(user_id: str): 206 """Возвращает список местоположений пользователя.""" 207 locations = thing_manager.get_all_things(user_id) 208 return {"locations": locations} 209 210 211@web_router.post("/api/locations") 212async def web_add_location(user_id: str, thing: str, place: str): 213 """Добавляет новую запись о местоположении вещи.""" 214 loc_id = thing_manager.add_location(user_id, thing, place) 215 return {"location_id": loc_id, "thing": thing, "place": place} 216 217 218@asynccontextmanager 219async def lifespan(app: FastAPI): 220 """ 221 Управляет жизненым циклом FastAPI-приложения. 222 223 При запуске настраивает логирование и запускает APScheduler 224 с фоновыми ассистенты (очистка Langfuse). 225 При остановке корректно завершает планировщик. 226 """ 227 setup_logging() 228 229 products_log = None 230 try: 231 from e_lib.product_log import ProductsLog 232 233 products_log = ProductsLog("LostAndFound", APP_VERSION, logger=logger) 234 products_log.post_started() 235 except Exception as e: 236 logger.warning(f"ProductsLog: {e}") 237 238 command_parser.cache_prompt() 239 240 scheduler = AsyncIOScheduler() 241 242 def to_cron(time_str: str) -> str: 243 """Преобразует время в формате "ЧЧ:ММ" в cron-выражение.""" 244 h, m = time_str.split(":") 245 return f"{m} {h} * * *" 246 247 assistants = config.get("Assistants", {}) 248 langfuse_cleaner_cfg = assistants.get("LangfuseCleaner", {}) 249 if langfuse_cleaner_cfg.get("ENABLED"): 250 lf_cleaner = LangfuseCleanerAssistant( 251 project_slug=langfuse_cleaner_cfg.get("PROJECT", ""), 252 period=langfuse_cleaner_cfg.get("RETENTION_DAYS", 7), 253 logger=logger, 254 ) 255 scheduler.add_job( 256 lf_cleaner.run, 257 CronTrigger.from_crontab(to_cron(langfuse_cleaner_cfg["TIME"])), 258 id="langfuse_cleaner", 259 replace_existing=True, 260 ) 261 logger.info(f"Ассистент '{langfuse_cleaner_cfg.get('NAME')}' включён") 262 263 llm_warmup_cfg = assistants.get("LLMWarmup", {}) 264 if llm_warmup_cfg.get("ENABLED"): 265 global llm_warmup_assistant 266 llm_warmup_assistant = LLMWarmupAssistant( 267 prompt=llm_warmup_cfg.get("PROMPT", "Привет"), 268 model=llm_warmup_cfg.get("MODEL", "e.anisimov/eagent"), 269 retries=llm_warmup_cfg.get("RETRIES", 3), 270 delay=llm_warmup_cfg.get("DELAY", 5), 271 logger=logger, 272 ) 273 logger.info(f"Ассистент '{llm_warmup_cfg.get('NAME')}' включён") 274 275 scheduler.start() 276 logger.info("Приложение запущено") 277 yield 278 scheduler.shutdown(wait=False) 279 logger.info("Приложение остановлено") 280 281 try: 282 if products_log: 283 products_log.post_shutdown() 284 except Exception as e: 285 logger.warning(f"ProductsLog: {e}") 286 287 288app = FastAPI( 289 title="Alice LostAndFound Skill", 290 description="Навык для Яндекс Алисы для учета вещей и их местоположений", 291 version="1.0.0", 292 lifespan=lifespan, 293) 294 295app.include_router(web_router, prefix="/web") 296 297if os.path.isdir(os.path.join(os.path.dirname(__file__), "web", "dist")): 298 app.mount( 299 "/web", 300 StaticFiles( 301 directory=os.path.join(os.path.dirname(__file__), "web", "dist"), html=True 302 ), 303 name="web", 304 ) 305 306 307@app.post("/") 308@observe(name="Alice LostAndFound Observability") 309async def handle_alice(request: AliceRequest): 310 """ 311 Основной обработчик запросов от Яндекс Алисы. 312 313 Принимает POST-запрос от Диалогов, проверяет skill_id, 314 парсит команду пользователя и выполняет соответствующее действие 315 (добавление, обновление, поиск вещей или вывод справки). 316 Возвращает ответ в формате Яндекс Диалогов. 317 318 Args: 319 request: Валидированный запрос от Яндекс Алиса 320 321 Returns: 322 dict: Ответ навыка в формате Яндекс Диалогов 323 324 Raises: 325 HTTPException: 403 если skill_id не совпадает 326 """ 327 request_id = uuid.uuid4().hex[:8] 328 log = logger.bind(request_id=request_id) 329 330 if request.session.skill_id != os.getenv("SKILL_ID"): 331 log.warning(f"Неверный skill_id: {request.session.skill_id}") 332 raise HTTPException(status_code=403, detail="Forbidden") 333 334 try: 335 user_id = request.session.user_id 336 # Маппим реальный user_id из сессии на логин владельца из config.yaml 337 users_config = config.get("USERS", {}) 338 for user in users_config: 339 if user_id in users_config[user]: 340 user_id = user 341 break 342 343 utterance = request.request.original_utterance.strip() 344 if not utterance: 345 log.info("Пустой запрос (старт сессии) — отвечаем пустой строкой") 346 _trigger_llm_warmup() 347 return build_response(request, "") 348 command = request.request.command.lower() 349 350 log.info(f"Запрос от пользователя {user_id}: {command}") 351 352 parsed = command_parser.parse(utterance) 353 response_tts = None 354 db_operation = "NONE" 355 thing = parsed.thing if parsed else None 356 place = parsed.place if parsed else None 357 358 if parsed is None: 359 response_text = "Команда не распознана" 360 response_tts = "-" 361 elif parsed.action == "ping": 362 response_text = "pong" 363 elif parsed.action in ("insert", "update"): 364 # Python code определяет UPDATE vs INSERT на основе существования записи в БД 365 existing = thing_manager.find_thing(user_id, parsed.thing) 366 if existing: 367 # Запись уже существует — обновляем место 368 if thing_manager.update_location(user_id, parsed.thing, parsed.place): 369 db_operation = "UPDATE" 370 response_text = ( 371 f"Место для '{parsed.thing}' обновлено на '{parsed.place}'." 372 ) 373 else: 374 loc_id = thing_manager.add_location( 375 user_id, parsed.thing, parsed.place 376 ) 377 db_operation = "INSERT" 378 response_text = f"Запись '{parsed.thing}' → '{parsed.place}' добавлена. ID: {loc_id}" 379 else: 380 # Новая запись 381 loc_id = thing_manager.add_location(user_id, parsed.thing, parsed.place) 382 db_operation = "INSERT" 383 response_text = f"Запись '{parsed.thing}' → '{parsed.place}' добавлена. ID: {loc_id}" 384 elif parsed.action == "query": 385 db_operation = "QUERY" 386 locations = thing_manager.find_thing(user_id, parsed.thing) 387 388 if not locations: 389 response_text = f"Записей о '{parsed.thing}' не найдено." 390 else: 391 places = [loc["place"] for loc in locations] 392 # Возвращаем последнее место (самое свежее) 393 if len(places) > 1: 394 response_text = ( 395 f"{parsed.thing}: {', '.join(places)}. Последнее — {places[-1]}" 396 ) 397 response_tts = ", ".join(places) 398 else: 399 response_text = f"{parsed.thing} {places[-1]}" 400 elif parsed.action == "search": 401 db_operation = "SEARCH" 402 results = thing_manager.search_thing(user_id, parsed.thing) 403 404 if not results: 405 response_text = f"По запросу '{parsed.thing}' ничего не найдено." 406 else: 407 lines = [f"{r['thing']} → {r['place']}" for r in results] 408 response_text = "Найдено: " + "; ".join(lines) 409 elif parsed.action == "delete": 410 db_operation = "DELETE" 411 if thing_manager.delete_location(user_id, parsed.thing): 412 response_text = f"Запись '{parsed.thing}' удалена." 413 else: 414 response_text = f"Записей о '{parsed.thing}' не найдено." 415 else: 416 response_text = "Команда не распознана" 417 response_tts = "-" 418 419 try: 420 get_langfuse().update_current_span( 421 metadata={ 422 "user_id": user_id, 423 "utterance": utterance, 424 "action": parsed.action if parsed else None, 425 "db_operation": db_operation, 426 "thing": thing, 427 "place": place, 428 "command": command, 429 "text_length": len(response_text), 430 "tts_length": len(response_tts) if response_tts else 0, 431 } 432 ) 433 get_langfuse().update_current_span(output=response_text) 434 except Exception as e: 435 log.warning(f"Langfuse: не удалось обновить span: {e}") 436 437 log.info(f"Отправляем ответ: {response_text}") 438 return build_response(request, response_text, tts=response_tts) 439 440 except Exception as e: 441 log.error(f"Ошибка обработки запроса: {e}") 442 return build_error_response() 443 444 445@app.get("/") 446async def root(): 447 """ 448 Корневой эндпоинт для проверки доступности навыка. 449 450 Returns: 451 dict: Статус приложения и путь к файлу данных 452 """ 453 return { 454 "message": "Alice LostAndFound Skill is running", 455 "status": "ok", 456 "data_file": DB_FILE, 457 } 458 459 460@app.get("/health") 461async def health_check(): 462 """ 463 Эндпоинт проверки здоровья приложения. 464 465 Returns: 466 dict: Статус health, временная метка и количество записей 467 """ 468 return { 469 "status": "healthy", 470 "timestamp": datetime.now().isoformat(), 471 "locations_count": thing_manager.get_location_count(), 472 } 473 474 475if __name__ == "__main__": 476 import uvicorn 477 478 logger.info("Запускам сервер на 0.0.0.0:8000") 479 uvicorn.run(app, host="0.0.0.0", port=8000) 480 481 482def _load_env(): 483 """ 484 Загружаем переменные из .env файла в os.environ. 485 486 Читает файлы .env и устанавливает: 487 - LANGFUSE_PUBLIC_KEY, LANGFUSE_SECRET_KEY, LANGFUSE_HOST 488 - LLM_BASE_URL, LLM_API_KEY 489 - SKILL_ID 490 """ 491 env_file = os.path.join(os.path.dirname(__file__), ".env") 492 if not os.path.exists(env_file): 493 logger.warning(f"Файл {env_file} не найден") 494 return 495 496 with open(env_file, "r", encoding="utf-8") as f: 497 for line in f: 498 line = line.strip() 499 if not line or line.startswith("#"): 500 continue 501 if "=" in line: 502 key, value = line.split("=", 1) 503 os.environ[key.strip()] = value.strip() 504 505 logger.info( 506 f"ENV загружен: LANGFUSE_PUBLIC_KEY={os.getenv('LANGFUSE_PUBLIC_KEY', '***')}, " 507 f"LLM_BASE_URL={os.getenv('LLM_BASE_URL', '***')}" 508 ) 509 510 511_load_env()
60def get_langfuse() -> Langfuse: 61 """ 62 Возвращает экземпляр Langfuse для observability (синглтон). 63 64 При первом вызове создаёт клиент с учётом переменных окружения 65 LANGFUSE_PUBLIC_KEY, LANGFUSE_SECRET_KEY и LANGFUSE_HOST. 66 67 Returns: 68 Langfuse: Экземпляр Langfuse-клиента 69 """ 70 global _langfuse_instance 71 if _langfuse_instance is None: 72 _langfuse_instance = Langfuse( 73 public_key=os.getenv("LANGFUSE_PUBLIC_KEY"), 74 secret_key=os.getenv("LANGFUSE_SECRET_KEY"), 75 host=os.getenv("LANGFUSE_HOST"), 76 ) 77 return _langfuse_instance
Возвращает экземпляр Langfuse для observability (синглтон).
При первом вызове создаёт клиент с учётом переменных окружения LANGFUSE_PUBLIC_KEY, LANGFUSE_SECRET_KEY и LANGFUSE_HOST.
Returns: Langfuse: Экземпляр Langfuse-клиента
83class InterceptHandler(logging.Handler): 84 """ 85 Перехватчик логов из стандартного logging в loguru. 86 87 Перенаправляет все сообщения из логгеров uvicorn и fastapi 88 в форматированный вывод loguru. 89 """ 90 91 def emit(self, record): 92 """ 93 Обрабатывает и перенаправляет запись лога в loguru. 94 95 Args: 96 record: Запись лога из стандартного модуля logging 97 """ 98 try: 99 level = logger.level(record.levelname).name 100 except ValueError: 101 level = record.levelno 102 103 frame, depth = logging.currentframe(), 2 104 while frame and frame.f_code.co_filename == logging.__file__: 105 frame = frame.f_back 106 depth += 1 107 108 logger.opt(depth=depth, exception=record.exc_info).log( 109 level, record.getMessage() 110 )
Перехватчик логов из стандартного logging в loguru.
Перенаправляет все сообщения из логгеров uvicorn и fastapi в форматированный вывод loguru.
91 def emit(self, record): 92 """ 93 Обрабатывает и перенаправляет запись лога в loguru. 94 95 Args: 96 record: Запись лога из стандартного модуля logging 97 """ 98 try: 99 level = logger.level(record.levelname).name 100 except ValueError: 101 level = record.levelno 102 103 frame, depth = logging.currentframe(), 2 104 while frame and frame.f_code.co_filename == logging.__file__: 105 frame = frame.f_back 106 depth += 1 107 108 logger.opt(depth=depth, exception=record.exc_info).log( 109 level, record.getMessage() 110 )
Обрабатывает и перенаправляет запись лога в loguru.
Args: record: Запись лога из стандартного модуля logging
113def setup_logging(): 114 """ 115 Настраивает логирование приложения. 116 117 Перенаправляет логи uvicorn и fastapi в loguru с цветным 118 форматированным выводом в stdout. Уровень DEBUG для всех 119 сообщений приложения. 120 """ 121 logger.remove() 122 123 logger.add( 124 sys.stdout, 125 format="<green>{time:YYYY-MM-DD HH:mm:ss}</green> | <level>{level: <8}</level> | <cyan>{name}</cyan>:<cyan>{function}</cyan>:<cyan>{line}</cyan> | <level>{message}</level>", 126 level="DEBUG", 127 colorize=True, 128 ) 129 130 intercept_handler = InterceptHandler() 131 132 root_logger = logging.getLogger() 133 root_logger.handlers = [intercept_handler] 134 root_logger.setLevel(logging.INFO) 135 136 for logger_name in ["uvicorn", "uvicorn.access", "uvicorn.error", "fastapi"]: 137 logging_logger = logging.getLogger(logger_name) 138 logging_logger.handlers = [intercept_handler] 139 logging_logger.setLevel(logging.INFO) 140 logging_logger.propagate = False
Настраивает логирование приложения.
Перенаправляет логи uvicorn и fastapi в loguru с цветным форматированным выводом в stdout. Уровень DEBUG для всех сообщений приложения.
143@observe(name="Build response") 144def build_response(request: AliceRequest, text: str, tts: str | None = None) -> dict: 145 """ 146 Формирует ответ навыка Яндекс Алисы. 147 148 Args: 149 request: Входящий запрос Алиса (для копирования session/version) 150 text: Текстовый ответ навыка 151 tts: Озвучивемый текст (если отличается от text) 152 153 Returns: 154 dict: Ответ в формате Яндекс Диалогов 155 """ 156 response: dict = { 157 "text": text, 158 "end_session": False, 159 } 160 if tts is not None: 161 response["tts"] = tts 162 return { 163 "version": request.version, 164 "session": { 165 "message_id": request.session.message_id, 166 "session_id": request.session.session_id, 167 "user_id": request.session.user_id, 168 "skill_id": request.session.skill_id, 169 }, 170 "response": response, 171 }
Формирует ответ навыка Яндекс Алисы.
Args: request: Входящий запрос Алиса (для копирования session/version) text: Текстовый ответ навыка tts: Озвучивемый текст (если отличается от text)
Returns: dict: Ответ в формате Яндекс Диалогов
174@observe(name="Build error response") 175def build_error_response( 176 text: str = "Произошла ошибка при обработке запроса. Попробуйте еще раз.", 177) -> dict: 178 """ 179 Формирует ответ с сообщением об ошибке. 180 181 Args: 182 text: Текст ошибки (по умолчанию стандартное сообщение) 183 184 Returns: 185 dict: Ответ-заглушка с текстом ошибки 186 """ 187 return { 188 "version": "1.0", 189 "session": { 190 "message_id": 0, 191 "session_id": "", 192 "user_id": "", 193 "skill_id": "", 194 }, 195 "response": { 196 "text": text, 197 "end_session": False, 198 }, 199 }
Формирует ответ с сообщением об ошибке.
Args: text: Текст ошибки (по умолчанию стандартное сообщение)
Returns: dict: Ответ-заглушка с текстом ошибки
205@web_router.get("/api/locations") 206async def web_list_locations(user_id: str): 207 """Возвращает список местоположений пользователя.""" 208 locations = thing_manager.get_all_things(user_id) 209 return {"locations": locations}
Возвращает список местоположений пользователя.
212@web_router.post("/api/locations") 213async def web_add_location(user_id: str, thing: str, place: str): 214 """Добавляет новую запись о местоположении вещи.""" 215 loc_id = thing_manager.add_location(user_id, thing, place) 216 return {"location_id": loc_id, "thing": thing, "place": place}
Добавляет новую запись о местоположении вещи.
219@asynccontextmanager 220async def lifespan(app: FastAPI): 221 """ 222 Управляет жизненым циклом FastAPI-приложения. 223 224 При запуске настраивает логирование и запускает APScheduler 225 с фоновыми ассистенты (очистка Langfuse). 226 При остановке корректно завершает планировщик. 227 """ 228 setup_logging() 229 230 products_log = None 231 try: 232 from e_lib.product_log import ProductsLog 233 234 products_log = ProductsLog("LostAndFound", APP_VERSION, logger=logger) 235 products_log.post_started() 236 except Exception as e: 237 logger.warning(f"ProductsLog: {e}") 238 239 command_parser.cache_prompt() 240 241 scheduler = AsyncIOScheduler() 242 243 def to_cron(time_str: str) -> str: 244 """Преобразует время в формате "ЧЧ:ММ" в cron-выражение.""" 245 h, m = time_str.split(":") 246 return f"{m} {h} * * *" 247 248 assistants = config.get("Assistants", {}) 249 langfuse_cleaner_cfg = assistants.get("LangfuseCleaner", {}) 250 if langfuse_cleaner_cfg.get("ENABLED"): 251 lf_cleaner = LangfuseCleanerAssistant( 252 project_slug=langfuse_cleaner_cfg.get("PROJECT", ""), 253 period=langfuse_cleaner_cfg.get("RETENTION_DAYS", 7), 254 logger=logger, 255 ) 256 scheduler.add_job( 257 lf_cleaner.run, 258 CronTrigger.from_crontab(to_cron(langfuse_cleaner_cfg["TIME"])), 259 id="langfuse_cleaner", 260 replace_existing=True, 261 ) 262 logger.info(f"Ассистент '{langfuse_cleaner_cfg.get('NAME')}' включён") 263 264 llm_warmup_cfg = assistants.get("LLMWarmup", {}) 265 if llm_warmup_cfg.get("ENABLED"): 266 global llm_warmup_assistant 267 llm_warmup_assistant = LLMWarmupAssistant( 268 prompt=llm_warmup_cfg.get("PROMPT", "Привет"), 269 model=llm_warmup_cfg.get("MODEL", "e.anisimov/eagent"), 270 retries=llm_warmup_cfg.get("RETRIES", 3), 271 delay=llm_warmup_cfg.get("DELAY", 5), 272 logger=logger, 273 ) 274 logger.info(f"Ассистент '{llm_warmup_cfg.get('NAME')}' включён") 275 276 scheduler.start() 277 logger.info("Приложение запущено") 278 yield 279 scheduler.shutdown(wait=False) 280 logger.info("Приложение остановлено") 281 282 try: 283 if products_log: 284 products_log.post_shutdown() 285 except Exception as e: 286 logger.warning(f"ProductsLog: {e}")
Управляет жизненым циклом FastAPI-приложения.
При запуске настраивает логирование и запускает APScheduler с фоновыми ассистенты (очистка Langfuse). При остановке корректно завершает планировщик.
308@app.post("/") 309@observe(name="Alice LostAndFound Observability") 310async def handle_alice(request: AliceRequest): 311 """ 312 Основной обработчик запросов от Яндекс Алисы. 313 314 Принимает POST-запрос от Диалогов, проверяет skill_id, 315 парсит команду пользователя и выполняет соответствующее действие 316 (добавление, обновление, поиск вещей или вывод справки). 317 Возвращает ответ в формате Яндекс Диалогов. 318 319 Args: 320 request: Валидированный запрос от Яндекс Алиса 321 322 Returns: 323 dict: Ответ навыка в формате Яндекс Диалогов 324 325 Raises: 326 HTTPException: 403 если skill_id не совпадает 327 """ 328 request_id = uuid.uuid4().hex[:8] 329 log = logger.bind(request_id=request_id) 330 331 if request.session.skill_id != os.getenv("SKILL_ID"): 332 log.warning(f"Неверный skill_id: {request.session.skill_id}") 333 raise HTTPException(status_code=403, detail="Forbidden") 334 335 try: 336 user_id = request.session.user_id 337 # Маппим реальный user_id из сессии на логин владельца из config.yaml 338 users_config = config.get("USERS", {}) 339 for user in users_config: 340 if user_id in users_config[user]: 341 user_id = user 342 break 343 344 utterance = request.request.original_utterance.strip() 345 if not utterance: 346 log.info("Пустой запрос (старт сессии) — отвечаем пустой строкой") 347 _trigger_llm_warmup() 348 return build_response(request, "") 349 command = request.request.command.lower() 350 351 log.info(f"Запрос от пользователя {user_id}: {command}") 352 353 parsed = command_parser.parse(utterance) 354 response_tts = None 355 db_operation = "NONE" 356 thing = parsed.thing if parsed else None 357 place = parsed.place if parsed else None 358 359 if parsed is None: 360 response_text = "Команда не распознана" 361 response_tts = "-" 362 elif parsed.action == "ping": 363 response_text = "pong" 364 elif parsed.action in ("insert", "update"): 365 # Python code определяет UPDATE vs INSERT на основе существования записи в БД 366 existing = thing_manager.find_thing(user_id, parsed.thing) 367 if existing: 368 # Запись уже существует — обновляем место 369 if thing_manager.update_location(user_id, parsed.thing, parsed.place): 370 db_operation = "UPDATE" 371 response_text = ( 372 f"Место для '{parsed.thing}' обновлено на '{parsed.place}'." 373 ) 374 else: 375 loc_id = thing_manager.add_location( 376 user_id, parsed.thing, parsed.place 377 ) 378 db_operation = "INSERT" 379 response_text = f"Запись '{parsed.thing}' → '{parsed.place}' добавлена. ID: {loc_id}" 380 else: 381 # Новая запись 382 loc_id = thing_manager.add_location(user_id, parsed.thing, parsed.place) 383 db_operation = "INSERT" 384 response_text = f"Запись '{parsed.thing}' → '{parsed.place}' добавлена. ID: {loc_id}" 385 elif parsed.action == "query": 386 db_operation = "QUERY" 387 locations = thing_manager.find_thing(user_id, parsed.thing) 388 389 if not locations: 390 response_text = f"Записей о '{parsed.thing}' не найдено." 391 else: 392 places = [loc["place"] for loc in locations] 393 # Возвращаем последнее место (самое свежее) 394 if len(places) > 1: 395 response_text = ( 396 f"{parsed.thing}: {', '.join(places)}. Последнее — {places[-1]}" 397 ) 398 response_tts = ", ".join(places) 399 else: 400 response_text = f"{parsed.thing} {places[-1]}" 401 elif parsed.action == "search": 402 db_operation = "SEARCH" 403 results = thing_manager.search_thing(user_id, parsed.thing) 404 405 if not results: 406 response_text = f"По запросу '{parsed.thing}' ничего не найдено." 407 else: 408 lines = [f"{r['thing']} → {r['place']}" for r in results] 409 response_text = "Найдено: " + "; ".join(lines) 410 elif parsed.action == "delete": 411 db_operation = "DELETE" 412 if thing_manager.delete_location(user_id, parsed.thing): 413 response_text = f"Запись '{parsed.thing}' удалена." 414 else: 415 response_text = f"Записей о '{parsed.thing}' не найдено." 416 else: 417 response_text = "Команда не распознана" 418 response_tts = "-" 419 420 try: 421 get_langfuse().update_current_span( 422 metadata={ 423 "user_id": user_id, 424 "utterance": utterance, 425 "action": parsed.action if parsed else None, 426 "db_operation": db_operation, 427 "thing": thing, 428 "place": place, 429 "command": command, 430 "text_length": len(response_text), 431 "tts_length": len(response_tts) if response_tts else 0, 432 } 433 ) 434 get_langfuse().update_current_span(output=response_text) 435 except Exception as e: 436 log.warning(f"Langfuse: не удалось обновить span: {e}") 437 438 log.info(f"Отправляем ответ: {response_text}") 439 return build_response(request, response_text, tts=response_tts) 440 441 except Exception as e: 442 log.error(f"Ошибка обработки запроса: {e}") 443 return build_error_response()
Основной обработчик запросов от Яндекс Алисы.
Принимает POST-запрос от Диалогов, проверяет skill_id, парсит команду пользователя и выполняет соответствующее действие (добавление, обновление, поиск вещей или вывод справки). Возвращает ответ в формате Яндекс Диалогов.
Args: request: Валидированный запрос от Яндекс Алиса
Returns: dict: Ответ навыка в формате Яндекс Диалогов
Raises: HTTPException: 403 если skill_id не совпадает
446@app.get("/") 447async def root(): 448 """ 449 Корневой эндпоинт для проверки доступности навыка. 450 451 Returns: 452 dict: Статус приложения и путь к файлу данных 453 """ 454 return { 455 "message": "Alice LostAndFound Skill is running", 456 "status": "ok", 457 "data_file": DB_FILE, 458 }
Корневой эндпоинт для проверки доступности навыка.
Returns: dict: Статус приложения и путь к файлу данных
461@app.get("/health") 462async def health_check(): 463 """ 464 Эндпоинт проверки здоровья приложения. 465 466 Returns: 467 dict: Статус health, временная метка и количество записей 468 """ 469 return { 470 "status": "healthy", 471 "timestamp": datetime.now().isoformat(), 472 "locations_count": thing_manager.get_location_count(), 473 }
Эндпоинт проверки здоровья приложения.
Returns: dict: Статус health, временная метка и количество записей