main
FastAPI-приложение навыка управления задачами для Яндекс Алисы.
Основная точка входа: обрабатывает POST-запросы от Яндекс Диалогов, парсит команды пользователя, выполняет CRUD-операции с задачами и возвращает ответы в формате JSON. Также предоставляет health-check и включает фоновые ассистенты (перенос задач, очистка Langfuse) через APScheduler.
1""" 2FastAPI-приложение навыка управления задачами для Яндекс Алисы. 3 4Основная точка входа: обрабатывает POST-запросы от Яндекс Диалогов, 5парсит команды пользователя, выполняет CRUD-операции с задачами 6и возвращает ответы в формате JSON. Также предоставляет health-check 7и включает фоновые ассистенты (перенос задач, очистка Langfuse) 8через APScheduler. 9""" 10 11import logging 12import os 13import sys 14import uuid 15import asyncio 16 17from fastapi import FastAPI, HTTPException, Query, APIRouter 18from fastapi.staticfiles import StaticFiles 19from datetime import datetime 20from typing import Optional 21from pydantic import BaseModel 22from fastapi.concurrency import asynccontextmanager 23from apscheduler.schedulers.asyncio import AsyncIOScheduler 24from apscheduler.triggers.cron import CronTrigger 25from langfuse import Langfuse 26from langfuse import observe 27from loguru import logger 28 29from e_lib.assistants import LangfuseCleanerAssistant, LLMWarmupAssistant 30from Assistants.otm import OTMAssistant 31from command_parser import CommandParser 32from e_lib import Config 33 34config = Config("config.yaml") 35APP_VERSION = os.getenv("PRODUCTS_VERSION", "dev") 36from models import AliceRequest 37from taskman import DB_FILE, TaskManager 38 39_langfuse_instance = None 40 41llm_warmup_assistant = None 42_llm_warmup_tasks = set() 43 44 45def _trigger_llm_warmup() -> None: 46 """ 47 Запускает фоновый прогрев LLM (best-effort). 48 49 Вызывается при старте сессии навыка: пока пользователь формулирует 50 вопрос, ассистент подгружает модель лёгким запросом. Повторные вызовы 51 при активном прогоне пропускаются самим ассистентом. Ошибки не роняют 52 запрос — прогрев строго фоновый. 53 """ 54 global llm_warmup_assistant 55 if llm_warmup_assistant is None: 56 return 57 task = asyncio.create_task(llm_warmup_assistant.run()) 58 _llm_warmup_tasks.add(task) 59 task.add_done_callback(_llm_warmup_tasks.discard) 60 61 62def get_langfuse() -> Langfuse: 63 """ 64 Возвращает экземпляр Langfuse для observability (синглтон). 65 66 При первом вызове создаёт подключение к Langfuse 67 с использованием переменных окружения LANGFUSE_PUBLIC_KEY, 68 LANGFUSE_SECRET_KEY и LANGFUSE_HOST. 69 70 Returns: 71 Langfuse: Экземпляр Langfuse-клиента 72 """ 73 global _langfuse_instance 74 if _langfuse_instance is None: 75 _langfuse_instance = Langfuse( 76 public_key=os.getenv("LANGFUSE_PUBLIC_KEY"), 77 secret_key=os.getenv("LANGFUSE_SECRET_KEY"), 78 host=os.getenv("LANGFUSE_HOST"), 79 ) 80 return _langfuse_instance 81 82 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.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 ) 111 112 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 141 142 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 } 172 173 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 } 200 201 202class TaskCreate(BaseModel): 203 """Модель для создания задачи через веб-API.""" 204 205 user_id: str 206 text: str 207 date: Optional[str] = None 208 209 210class TaskUpdate(BaseModel): 211 """Модель для обновления задачи через веб-API.""" 212 213 date: Optional[str] = None 214 completed: Optional[bool] = None 215 216 217web_router = APIRouter() 218 219 220@web_router.get("/api/tasks") 221async def web_list_tasks( 222 user_id: str, 223 date: Optional[str] = None, 224 order: Optional[str] = Query( 225 None 226 ), # "newest-first" или "oldest-first", по умолч. "newest-first" 227): 228 """Возвращает список задач пользователя с фильтрацией по дате и сортировкой.""" 229 if date == "all": 230 tasks = task_manager.get_all_tasks(user_id) 231 232 # Сортировка: сначала новые (с датой), потом без даты, затем по убыванию даты 233 def sort_key(task): 234 if not task.get("date"): 235 return ("0", datetime.min) # Без даты - в конец 236 has_future_date = task["date"] != "2100-01-01" 237 if has_future_date: 238 # Реальная дата - сортируем по убыванию (новые первыми) 239 try: 240 dt = datetime.strptime(task["date"], "%Y-%m-%d") 241 return ("1", dt) 242 except ValueError: 243 return ("1", datetime.min) 244 else: 245 # Задачи без даты (2100-01-01) - в конец 246 return ("0", datetime.max) 247 248 if order == "oldest-first": 249 tasks.sort(key=sort_key, reverse=False) # oldest-first: ascending 250 else: 251 tasks.sort(key=sort_key, reverse=True) # newest-first: descending 252 253 return {"tasks": tasks, "date": "all"} 254 if date is None: 255 date = task_manager.get_today_date() 256 tasks = task_manager.get_tasks_by_date(user_id, date) 257 return {"tasks": tasks, "date": date} 258 259 260@web_router.post("/api/tasks") 261async def web_create_task(body: TaskCreate): 262 """Создаёт новую задачу.""" 263 task_id = task_manager.add_task(body.user_id, body.text, body.date) 264 task = task_manager.get_task(body.user_id, task_id) 265 return {"task": task} 266 267 268@web_router.put("/api/tasks/{task_id}") 269async def web_update_task(task_id: int, body: TaskUpdate, user_id: str = Query(...)): 270 """Обновляет задачу (дату и/или статус выполнения).""" 271 if body.date is not None: 272 from datetime import datetime as _dt 273 274 try: 275 date = _dt.strptime(body.date, "%Y-%m-%d").strftime("%Y-%m-%d") 276 except ValueError: 277 raise HTTPException( 278 status_code=400, detail="Invalid date format, expected YYYY-MM-DD" 279 ) 280 with task_manager._get_conn() as conn: 281 conn.execute( 282 "UPDATE tasks SET date = ? WHERE user_id = ? AND id = ?", 283 (date, user_id, task_id), 284 ) 285 if body.completed is not None: 286 if body.completed: 287 task_manager.complete_task(user_id, task_id) 288 task = task_manager.get_task(user_id, task_id) 289 if task is None: 290 raise HTTPException(status_code=404, detail="Task not found") 291 return {"task": task} 292 293 294@web_router.delete("/api/tasks/{task_id}") 295async def web_delete_task(task_id: int, user_id: str = Query(...)): 296 """Удаляет задачу по ID.""" 297 if task_manager.delete_task(user_id, task_id): 298 return {"deleted": True} 299 raise HTTPException(status_code=404, detail="Task not found") 300 301 302@web_router.get("/api/users") 303async def web_list_users(): 304 """Возвращает список пользователей из config.data.""" 305 users = {} 306 for name, ids in config.data.get("USERS", {}).items(): 307 users[name] = ids 308 return {"users": users} 309 310 311@asynccontextmanager 312async def lifespan(app: FastAPI): 313 """ 314 Управляет жизненным циклом FastAPI-приложения. 315 316 При запуске настраивает логирование и запускает APScheduler 317 с фоновыми ассистентами (перенос задач, очистка Langfuse). 318 При остановке корректно завершает планировщик. 319 """ 320 setup_logging() 321 322 try: 323 command_parser.cache_prompt() 324 except Exception as e: 325 logger.warning(f"Langfuse: ошибка инициализации промпта: {e}") 326 327 products_log = None 328 try: 329 from e_lib.product_log import ProductsLog 330 331 products_log = ProductsLog("AliceTodoSkill", APP_VERSION, logger=logger) 332 products_log.post_started() 333 except Exception as e: 334 products_log = None 335 logger.warning(f"ProductsLog: {e}") 336 337 scheduler = AsyncIOScheduler() 338 339 def to_cron(time_str: str) -> str: 340 """ 341 Преобразует время в формате HH:MM в cron-выражение. 342 343 Args: 344 time_str: Время в формате "ЧЧ:ММ" 345 346 Returns: 347 str: Cron-выражение для ежедневного запуска 348 """ 349 h, m = time_str.split(":") 350 return f"{m} {h} * * *" 351 352 if config.data["Assistants"]["TaskMover"]["ENABLED"]: 353 old_task_move_assistant = OTMAssistant( 354 name=config.data["Assistants"]["TaskMover"]["NAME"], 355 task_manager=task_manager, 356 ) 357 scheduler.add_job( 358 old_task_move_assistant.run, 359 CronTrigger.from_crontab( 360 to_cron(config.data["Assistants"]["TaskMover"]["TIME"]) 361 ), 362 id="old_task_mover", 363 replace_existing=True, 364 ) 365 logger.info( 366 f'Ассистент "{config.data["Assistants"]["TaskMover"]["NAME"]}" включён' 367 ) 368 369 if config.data["Assistants"]["LangfuseCleaner"]["ENABLED"]: 370 lf_cleaner_assistant = LangfuseCleanerAssistant( 371 project_slug=config.data["Assistants"]["LangfuseCleaner"]["PROJECT"], 372 period=config.data["Assistants"]["LangfuseCleaner"]["RETENTION_DAYS"], 373 logger=logger, 374 ) 375 scheduler.add_job( 376 lf_cleaner_assistant.run, 377 CronTrigger.from_crontab( 378 to_cron(config.data["Assistants"]["LangfuseCleaner"]["TIME"]) 379 ), 380 id="lf_cleaner", 381 replace_existing=True, 382 ) 383 logger.info( 384 f'Ассистент "{config.data["Assistants"]["LangfuseCleaner"]["NAME"]}" включён' 385 ) 386 387 if config.data["Assistants"]["LLMWarmup"]["ENABLED"]: 388 global llm_warmup_assistant 389 llm_warmup_assistant = LLMWarmupAssistant( 390 prompt=config.data["Assistants"]["LLMWarmup"]["PROMPT"], 391 model=config.data["Assistants"]["LLMWarmup"]["MODEL"], 392 retries=config.data["Assistants"]["LLMWarmup"]["RETRIES"], 393 delay=config.data["Assistants"]["LLMWarmup"]["DELAY"], 394 logger=logger, 395 ) 396 logger.info( 397 f'Ассистент "{config.data["Assistants"]["LLMWarmup"]["NAME"]}" включён' 398 ) 399 400 scheduler.start() 401 logger.info("Приложение запущено") 402 yield 403 scheduler.shutdown(wait=False) 404 logger.info("Приложение остановлено") 405 try: 406 if products_log: 407 products_log.post_shutdown() 408 except Exception as e: 409 logger.warning(f"ProductsLog: {e}") 410 411 412app = FastAPI( 413 title="Alice Todo Skill", 414 description="Навык для Яндекс Алисы для управления задачами", 415 version="1.0.0", 416 lifespan=lifespan, 417) 418 419task_manager = TaskManager() 420command_parser = CommandParser() 421 422app.include_router(web_router, prefix="/web") 423 424if os.path.isdir(os.path.join(os.path.dirname(__file__), "web", "dist")): 425 app.mount( 426 "/web", 427 StaticFiles( 428 directory=os.path.join(os.path.dirname(__file__), "web", "dist"), html=True 429 ), 430 name="web", 431 ) 432 433 434@app.post("/") 435@observe(name="Alice Todo Skill Observability") 436async def handle_alice(request: AliceRequest): 437 """ 438 Основной обработчик запросов от Яндекс Алисы. 439 440 Принимает POST-запрос от Диалогов, проверяет skill_id, 441 парсит команду пользователя и выполняет соответствующее действие 442 (добавление, просмотр, выполнение, удаление, перенос задач 443 или вывод справки). Возвращает ответ в формате Яндекс Диалогов. 444 445 Args: 446 request: Валидированный запрос от Яндекс Алисы 447 448 Returns: 449 dict: Ответ навыка в формате Яндекс Диалогов 450 451 Raises: 452 HTTPException: 403 если skill_id не совпадает 453 """ 454 request_id = uuid.uuid4().hex[:8] 455 log = logger.bind(request_id=request_id) 456 457 if request.session.skill_id != os.getenv("SKILL_ID"): 458 log.warning(f"Неверный skill_id: {request.session.skill_id}") 459 raise HTTPException(status_code=403, detail="Forbidden") 460 461 try: 462 user_id = request.session.user_id 463 for user in config.data["USERS"]: 464 if user_id in config.data["USERS"][user]: 465 user_id = user 466 break 467 468 command = request.request.command.lower() 469 original_utterance = request.request.original_utterance.lower() 470 471 log.info(f"Запрос от пользователя {user_id}: {command}") 472 473 if not original_utterance.strip(): 474 _trigger_llm_warmup() 475 return build_response(request, "") 476 477 parsed = command_parser.parse(command, original_utterance) 478 response_tts = None 479 480 if parsed is None: 481 response_text = "Мой ежедневник" 482 response_tts = "-" 483 elif parsed.action == "ping": 484 response_text = "pong" 485 elif parsed.action == "add_task": 486 if parsed.date_keyword is None: 487 task_id = task_manager.add_task(user_id, parsed.task_text) 488 response_text = ( 489 f"Задача '{parsed.task_text}' добавлена. ID задачи: {task_id}" 490 ) 491 else: 492 date = parsed.date_keyword 493 task_id = task_manager.add_task(user_id, parsed.task_text, date) 494 response_text = f"Задача '{parsed.task_text}' добавлена на {parsed.date_label or parsed.date_keyword}. ID задачи: {task_id}" 495 elif parsed.action == "no_task_text": 496 response_text = "Не смогла понять, какую задачу добавить." 497 elif parsed.action == "list_tasks": 498 if parsed.date_keyword is None: 499 tasks = task_manager.get_all_tasks(user_id) 500 incomplete_count = len(tasks) 501 scope = "all" 502 else: 503 date = parsed.date_keyword 504 tasks = task_manager.get_tasks_by_date(user_id, date) 505 incomplete_count = task_manager.get_incomplete_tasks_count( 506 user_id, date 507 ) 508 scope = "date" 509 510 if incomplete_count == 0: 511 if scope == "all": 512 response_text = "Невыполненных задач нет." 513 else: 514 response_text = f"На {parsed.date_label or parsed.date_keyword} невыполненных задач нет." 515 else: 516 if scope == "all": 517 prefix = f"У вас {incomplete_count} невыполненных задач: " 518 else: 519 prefix = f"На {parsed.date_label or parsed.date_keyword} у вас {incomplete_count} невыполненных задач: " 520 response_text = prefix + task_manager.format_task_list_truncated( 521 tasks, show_completed=False, max_chars=1024 - len(prefix) 522 ) 523 tts_formatted = task_manager.format_task_list_truncated( 524 tasks, show_completed=False, max_chars=1024 525 ) 526 if "\nи еще" in tts_formatted or tts_formatted.startswith( 527 "Не показано" 528 ): 529 summary = task_manager.get_incomplete_tasks_summary(user_id) 530 response_tts = ( 531 f"У вас {summary['total']} невыполненных задач. " 532 f"На сегодня {summary['today']}, на завтра {summary['tomorrow']}." 533 ) 534 else: 535 response_tts = tts_formatted 536 elif parsed.action == "help": 537 response_text = ( 538 "Я помогу вам управлять задачами. Вот что я умею:\n" 539 "• 'Добавь задачу купить молоко' - добавить задачу без срока\n" 540 "• 'Добавь задачу на завтра сходить к врачу' - добавить на завтра\n" 541 "• 'Добавь задачу на послезавтра ...' - добавить на послезавтра\n" 542 "• 'Какие задачи на сегодня?' - показать невыполненные задачи\n" 543 "• 'Дела на завтра' - показать задачи на завтра\n" 544 "• 'Что купить?' - найти задачи, начинающиеся с 'купить'\n" 545 "• 'Пометь задачу 3 выполненной' - отметить задачу как выполненную\n" 546 "• 'Удали задачу 2' / 'убери задачу 2' - удалить задачу\n" 547 "• 'Перенеси задачу 1 на завтра' - перенести одну задачу\n" 548 "• 'Перенеси задачи на завтра' - перенести все задачи" 549 ) 550 elif parsed.action == "what": 551 date = parsed.date_keyword 552 tasks = task_manager.get_tasks_by_date(user_id, date) 553 filtered_tasks = [ 554 t for t in tasks if t["text"].lower().startswith(parsed.keyword) 555 ] 556 incomplete_count = sum(1 for t in filtered_tasks if not t["completed"]) 557 558 if incomplete_count == 0: 559 response_text = f"На {parsed.date_label or parsed.date_keyword} нет задач, начинающихся с '{parsed.keyword}'." 560 else: 561 prefix = f"На {parsed.date_label or parsed.date_keyword} у вас {incomplete_count} задач: " 562 response_text = prefix + task_manager.format_task_list_truncated( 563 filtered_tasks, show_completed=False, max_chars=1024 - len(prefix) 564 ) 565 tts_formatted = task_manager.format_task_list_truncated( 566 filtered_tasks, show_completed=False, max_chars=1024 567 ) 568 if "\nи еще" in tts_formatted or tts_formatted.startswith( 569 "Не показано" 570 ): 571 summary = task_manager.get_incomplete_tasks_summary(user_id) 572 response_tts = ( 573 f"У вас {summary['total']} невыполненных задач. " 574 f"На сегодня {summary['today']}, на завтра {summary['tomorrow']}." 575 ) 576 else: 577 response_tts = tts_formatted 578 elif parsed.action == "complete": 579 if task_manager.complete_task(user_id, parsed.task_id): 580 response_text = f"Задача {parsed.task_id} отмечена как выполненная." 581 else: 582 response_text = f"Задача с номером {parsed.task_id} не найдена." 583 elif parsed.action == "delete": 584 if task_manager.delete_task(user_id, parsed.task_id): 585 response_text = f"Задача {parsed.task_id} удалена." 586 else: 587 response_text = f"Задача с номером {parsed.task_id} не найдена." 588 elif parsed.action == "no_task_id": 589 response_text = "Не смогла найти номер задачи." 590 elif parsed.action == "move_one": 591 date = parsed.date_keyword 592 if task_manager.move_task(user_id, parsed.task_id, date): 593 response_text = f"Задача {parsed.task_id} перенесена на {parsed.date_label or parsed.date_keyword}." 594 else: 595 response_text = f"Задача с номером {parsed.task_id} не найдена." 596 elif parsed.action == "move_all": 597 date = parsed.date_keyword 598 task_manager.move_tasks(user_id, date) 599 response_text = ( 600 f"Задачи перенесены на {parsed.date_label or parsed.date_keyword}." 601 ) 602 603 try: 604 get_langfuse().update_current_span( 605 metadata={ 606 "user_id": user_id, 607 "action": parsed.action if parsed else None, 608 "command": command, 609 "text_length": len(response_text), 610 "tts_length": len(response_tts) if response_tts else 0, 611 } 612 ) 613 except Exception: 614 pass 615 log.info(f"Отправляем ответ: {response_text}") 616 return build_response(request, response_text, tts=response_tts) 617 618 except Exception as e: 619 log.error(f"Ошибка обработки запроса: {e}") 620 return build_error_response() 621 622 623@app.get("/") 624async def root(): 625 """ 626 Корневой эндпоинт для проверки доступности навыка. 627 628 Returns: 629 dict: Статус приложения и путь к файлу данных 630 """ 631 return { 632 "message": "Alice Todo Skill is running", 633 "status": "ok", 634 "data_file": DB_FILE, 635 } 636 637 638@app.get("/health") 639async def health_check(): 640 """ 641 Эндпоинт проверки здоровья приложения. 642 643 Returns: 644 dict: Статус health, временная метка и количество пользователей 645 """ 646 return { 647 "status": "healthy", 648 "timestamp": datetime.now().isoformat(), 649 "users_count": task_manager.get_user_count(), 650 } 651 652 653if __name__ == "__main__": 654 import uvicorn 655 656 uvicorn.run(app, host="0.0.0.0", port=8000)
63def get_langfuse() -> Langfuse: 64 """ 65 Возвращает экземпляр Langfuse для observability (синглтон). 66 67 При первом вызове создаёт подключение к Langfuse 68 с использованием переменных окружения LANGFUSE_PUBLIC_KEY, 69 LANGFUSE_SECRET_KEY и LANGFUSE_HOST. 70 71 Returns: 72 Langfuse: Экземпляр Langfuse-клиента 73 """ 74 global _langfuse_instance 75 if _langfuse_instance is None: 76 _langfuse_instance = Langfuse( 77 public_key=os.getenv("LANGFUSE_PUBLIC_KEY"), 78 secret_key=os.getenv("LANGFUSE_SECRET_KEY"), 79 host=os.getenv("LANGFUSE_HOST"), 80 ) 81 return _langfuse_instance
Возвращает экземпляр Langfuse для observability (синглтон).
При первом вызове создаёт подключение к Langfuse с использованием переменных окружения LANGFUSE_PUBLIC_KEY, LANGFUSE_SECRET_KEY и LANGFUSE_HOST.
Returns: Langfuse: Экземпляр Langfuse-клиента
84class InterceptHandler(logging.Handler): 85 """ 86 Перехватчик логов из стандартного logging в loguru. 87 88 Перенаправляет все сообщения из логгеров uvicorn и fastapi 89 в форматированный вывод loguru. 90 """ 91 92 def emit(self, record): 93 """ 94 Обрабатывает и перенаправляет запись лога в loguru. 95 96 Args: 97 record: Запись лога из стандартного модуля logging 98 """ 99 try: 100 level = logger.level(record.levelname).name 101 except ValueError: 102 level = record.levelno 103 104 frame, depth = logging.currentframe(), 2 105 while frame.f_code.co_filename == logging.__file__: 106 frame = frame.f_back 107 depth += 1 108 109 logger.opt(depth=depth, exception=record.exc_info).log( 110 level, record.getMessage() 111 )
Перехватчик логов из стандартного logging в loguru.
Перенаправляет все сообщения из логгеров uvicorn и fastapi в форматированный вывод loguru.
92 def emit(self, record): 93 """ 94 Обрабатывает и перенаправляет запись лога в loguru. 95 96 Args: 97 record: Запись лога из стандартного модуля logging 98 """ 99 try: 100 level = logger.level(record.levelname).name 101 except ValueError: 102 level = record.levelno 103 104 frame, depth = logging.currentframe(), 2 105 while frame.f_code.co_filename == logging.__file__: 106 frame = frame.f_back 107 depth += 1 108 109 logger.opt(depth=depth, exception=record.exc_info).log( 110 level, record.getMessage() 111 )
Обрабатывает и перенаправляет запись лога в loguru.
Args: record: Запись лога из стандартного модуля logging
114def setup_logging(): 115 """ 116 Настраивает логирование приложения. 117 118 Перенаправляет логи uvicorn и fastapi в loguru с цветным 119 форматированным выводом в stdout. Уровень DEBUG для всех 120 сообщений приложения. 121 """ 122 logger.remove() 123 124 logger.add( 125 sys.stdout, 126 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>", 127 level="DEBUG", 128 colorize=True, 129 ) 130 131 intercept_handler = InterceptHandler() 132 133 root_logger = logging.getLogger() 134 root_logger.handlers = [intercept_handler] 135 root_logger.setLevel(logging.INFO) 136 137 for logger_name in ["uvicorn", "uvicorn.access", "uvicorn.error", "fastapi"]: 138 logging_logger = logging.getLogger(logger_name) 139 logging_logger.handlers = [intercept_handler] 140 logging_logger.setLevel(logging.INFO) 141 logging_logger.propagate = False
Настраивает логирование приложения.
Перенаправляет логи uvicorn и fastapi в loguru с цветным форматированным выводом в stdout. Уровень DEBUG для всех сообщений приложения.
144@observe(name="Build response") 145def build_response(request: AliceRequest, text: str, tts: str | None = None) -> dict: 146 """ 147 Формирует ответ навыка Яндекс Алисы. 148 149 Args: 150 request: Входящий запрос Алисы (для копирования session/version) 151 text: Текстовый ответ навыка 152 tts: Озвучиваемый текст (если отличается от text) 153 154 Returns: 155 dict: Ответ в формате Яндекс Диалогов 156 """ 157 response: dict = { 158 "text": text, 159 "end_session": False, 160 } 161 if tts is not None: 162 response["tts"] = tts 163 return { 164 "version": request.version, 165 "session": { 166 "message_id": request.session.message_id, 167 "session_id": request.session.session_id, 168 "user_id": request.session.user_id, 169 "skill_id": request.session.skill_id, 170 }, 171 "response": response, 172 }
Формирует ответ навыка Яндекс Алисы.
Args: request: Входящий запрос Алисы (для копирования session/version) text: Текстовый ответ навыка tts: Озвучиваемый текст (если отличается от text)
Returns: dict: Ответ в формате Яндекс Диалогов
175@observe(name="Build error response") 176def build_error_response( 177 text: str = "Произошла ошибка при обработке запроса. Попробуйте еще раз.", 178) -> dict: 179 """ 180 Формирует ответ с сообщением об ошибке. 181 182 Args: 183 text: Текст ошибки (по умолчанию стандартное сообщение) 184 185 Returns: 186 dict: Ответ-заглушка с текстом ошибки 187 """ 188 return { 189 "version": "1.0", 190 "session": { 191 "message_id": 0, 192 "session_id": "", 193 "user_id": "", 194 "skill_id": "", 195 }, 196 "response": { 197 "text": text, 198 "end_session": False, 199 }, 200 }
Формирует ответ с сообщением об ошибке.
Args: text: Текст ошибки (по умолчанию стандартное сообщение)
Returns: dict: Ответ-заглушка с текстом ошибки
203class TaskCreate(BaseModel): 204 """Модель для создания задачи через веб-API.""" 205 206 user_id: str 207 text: str 208 date: Optional[str] = None
Модель для создания задачи через веб-API.
211class TaskUpdate(BaseModel): 212 """Модель для обновления задачи через веб-API.""" 213 214 date: Optional[str] = None 215 completed: Optional[bool] = None
Модель для обновления задачи через веб-API.
221@web_router.get("/api/tasks") 222async def web_list_tasks( 223 user_id: str, 224 date: Optional[str] = None, 225 order: Optional[str] = Query( 226 None 227 ), # "newest-first" или "oldest-first", по умолч. "newest-first" 228): 229 """Возвращает список задач пользователя с фильтрацией по дате и сортировкой.""" 230 if date == "all": 231 tasks = task_manager.get_all_tasks(user_id) 232 233 # Сортировка: сначала новые (с датой), потом без даты, затем по убыванию даты 234 def sort_key(task): 235 if not task.get("date"): 236 return ("0", datetime.min) # Без даты - в конец 237 has_future_date = task["date"] != "2100-01-01" 238 if has_future_date: 239 # Реальная дата - сортируем по убыванию (новые первыми) 240 try: 241 dt = datetime.strptime(task["date"], "%Y-%m-%d") 242 return ("1", dt) 243 except ValueError: 244 return ("1", datetime.min) 245 else: 246 # Задачи без даты (2100-01-01) - в конец 247 return ("0", datetime.max) 248 249 if order == "oldest-first": 250 tasks.sort(key=sort_key, reverse=False) # oldest-first: ascending 251 else: 252 tasks.sort(key=sort_key, reverse=True) # newest-first: descending 253 254 return {"tasks": tasks, "date": "all"} 255 if date is None: 256 date = task_manager.get_today_date() 257 tasks = task_manager.get_tasks_by_date(user_id, date) 258 return {"tasks": tasks, "date": date}
Возвращает список задач пользователя с фильтрацией по дате и сортировкой.
261@web_router.post("/api/tasks") 262async def web_create_task(body: TaskCreate): 263 """Создаёт новую задачу.""" 264 task_id = task_manager.add_task(body.user_id, body.text, body.date) 265 task = task_manager.get_task(body.user_id, task_id) 266 return {"task": task}
Создаёт новую задачу.
269@web_router.put("/api/tasks/{task_id}") 270async def web_update_task(task_id: int, body: TaskUpdate, user_id: str = Query(...)): 271 """Обновляет задачу (дату и/или статус выполнения).""" 272 if body.date is not None: 273 from datetime import datetime as _dt 274 275 try: 276 date = _dt.strptime(body.date, "%Y-%m-%d").strftime("%Y-%m-%d") 277 except ValueError: 278 raise HTTPException( 279 status_code=400, detail="Invalid date format, expected YYYY-MM-DD" 280 ) 281 with task_manager._get_conn() as conn: 282 conn.execute( 283 "UPDATE tasks SET date = ? WHERE user_id = ? AND id = ?", 284 (date, user_id, task_id), 285 ) 286 if body.completed is not None: 287 if body.completed: 288 task_manager.complete_task(user_id, task_id) 289 task = task_manager.get_task(user_id, task_id) 290 if task is None: 291 raise HTTPException(status_code=404, detail="Task not found") 292 return {"task": task}
Обновляет задачу (дату и/или статус выполнения).
295@web_router.delete("/api/tasks/{task_id}") 296async def web_delete_task(task_id: int, user_id: str = Query(...)): 297 """Удаляет задачу по ID.""" 298 if task_manager.delete_task(user_id, task_id): 299 return {"deleted": True} 300 raise HTTPException(status_code=404, detail="Task not found")
Удаляет задачу по ID.
303@web_router.get("/api/users") 304async def web_list_users(): 305 """Возвращает список пользователей из config.data.""" 306 users = {} 307 for name, ids in config.data.get("USERS", {}).items(): 308 users[name] = ids 309 return {"users": users}
Возвращает список пользователей из config.data.
312@asynccontextmanager 313async def lifespan(app: FastAPI): 314 """ 315 Управляет жизненным циклом FastAPI-приложения. 316 317 При запуске настраивает логирование и запускает APScheduler 318 с фоновыми ассистентами (перенос задач, очистка Langfuse). 319 При остановке корректно завершает планировщик. 320 """ 321 setup_logging() 322 323 try: 324 command_parser.cache_prompt() 325 except Exception as e: 326 logger.warning(f"Langfuse: ошибка инициализации промпта: {e}") 327 328 products_log = None 329 try: 330 from e_lib.product_log import ProductsLog 331 332 products_log = ProductsLog("AliceTodoSkill", APP_VERSION, logger=logger) 333 products_log.post_started() 334 except Exception as e: 335 products_log = None 336 logger.warning(f"ProductsLog: {e}") 337 338 scheduler = AsyncIOScheduler() 339 340 def to_cron(time_str: str) -> str: 341 """ 342 Преобразует время в формате HH:MM в cron-выражение. 343 344 Args: 345 time_str: Время в формате "ЧЧ:ММ" 346 347 Returns: 348 str: Cron-выражение для ежедневного запуска 349 """ 350 h, m = time_str.split(":") 351 return f"{m} {h} * * *" 352 353 if config.data["Assistants"]["TaskMover"]["ENABLED"]: 354 old_task_move_assistant = OTMAssistant( 355 name=config.data["Assistants"]["TaskMover"]["NAME"], 356 task_manager=task_manager, 357 ) 358 scheduler.add_job( 359 old_task_move_assistant.run, 360 CronTrigger.from_crontab( 361 to_cron(config.data["Assistants"]["TaskMover"]["TIME"]) 362 ), 363 id="old_task_mover", 364 replace_existing=True, 365 ) 366 logger.info( 367 f'Ассистент "{config.data["Assistants"]["TaskMover"]["NAME"]}" включён' 368 ) 369 370 if config.data["Assistants"]["LangfuseCleaner"]["ENABLED"]: 371 lf_cleaner_assistant = LangfuseCleanerAssistant( 372 project_slug=config.data["Assistants"]["LangfuseCleaner"]["PROJECT"], 373 period=config.data["Assistants"]["LangfuseCleaner"]["RETENTION_DAYS"], 374 logger=logger, 375 ) 376 scheduler.add_job( 377 lf_cleaner_assistant.run, 378 CronTrigger.from_crontab( 379 to_cron(config.data["Assistants"]["LangfuseCleaner"]["TIME"]) 380 ), 381 id="lf_cleaner", 382 replace_existing=True, 383 ) 384 logger.info( 385 f'Ассистент "{config.data["Assistants"]["LangfuseCleaner"]["NAME"]}" включён' 386 ) 387 388 if config.data["Assistants"]["LLMWarmup"]["ENABLED"]: 389 global llm_warmup_assistant 390 llm_warmup_assistant = LLMWarmupAssistant( 391 prompt=config.data["Assistants"]["LLMWarmup"]["PROMPT"], 392 model=config.data["Assistants"]["LLMWarmup"]["MODEL"], 393 retries=config.data["Assistants"]["LLMWarmup"]["RETRIES"], 394 delay=config.data["Assistants"]["LLMWarmup"]["DELAY"], 395 logger=logger, 396 ) 397 logger.info( 398 f'Ассистент "{config.data["Assistants"]["LLMWarmup"]["NAME"]}" включён' 399 ) 400 401 scheduler.start() 402 logger.info("Приложение запущено") 403 yield 404 scheduler.shutdown(wait=False) 405 logger.info("Приложение остановлено") 406 try: 407 if products_log: 408 products_log.post_shutdown() 409 except Exception as e: 410 logger.warning(f"ProductsLog: {e}")
Управляет жизненным циклом FastAPI-приложения.
При запуске настраивает логирование и запускает APScheduler с фоновыми ассистентами (перенос задач, очистка Langfuse). При остановке корректно завершает планировщик.
435@app.post("/") 436@observe(name="Alice Todo Skill Observability") 437async def handle_alice(request: AliceRequest): 438 """ 439 Основной обработчик запросов от Яндекс Алисы. 440 441 Принимает POST-запрос от Диалогов, проверяет skill_id, 442 парсит команду пользователя и выполняет соответствующее действие 443 (добавление, просмотр, выполнение, удаление, перенос задач 444 или вывод справки). Возвращает ответ в формате Яндекс Диалогов. 445 446 Args: 447 request: Валидированный запрос от Яндекс Алисы 448 449 Returns: 450 dict: Ответ навыка в формате Яндекс Диалогов 451 452 Raises: 453 HTTPException: 403 если skill_id не совпадает 454 """ 455 request_id = uuid.uuid4().hex[:8] 456 log = logger.bind(request_id=request_id) 457 458 if request.session.skill_id != os.getenv("SKILL_ID"): 459 log.warning(f"Неверный skill_id: {request.session.skill_id}") 460 raise HTTPException(status_code=403, detail="Forbidden") 461 462 try: 463 user_id = request.session.user_id 464 for user in config.data["USERS"]: 465 if user_id in config.data["USERS"][user]: 466 user_id = user 467 break 468 469 command = request.request.command.lower() 470 original_utterance = request.request.original_utterance.lower() 471 472 log.info(f"Запрос от пользователя {user_id}: {command}") 473 474 if not original_utterance.strip(): 475 _trigger_llm_warmup() 476 return build_response(request, "") 477 478 parsed = command_parser.parse(command, original_utterance) 479 response_tts = None 480 481 if parsed is None: 482 response_text = "Мой ежедневник" 483 response_tts = "-" 484 elif parsed.action == "ping": 485 response_text = "pong" 486 elif parsed.action == "add_task": 487 if parsed.date_keyword is None: 488 task_id = task_manager.add_task(user_id, parsed.task_text) 489 response_text = ( 490 f"Задача '{parsed.task_text}' добавлена. ID задачи: {task_id}" 491 ) 492 else: 493 date = parsed.date_keyword 494 task_id = task_manager.add_task(user_id, parsed.task_text, date) 495 response_text = f"Задача '{parsed.task_text}' добавлена на {parsed.date_label or parsed.date_keyword}. ID задачи: {task_id}" 496 elif parsed.action == "no_task_text": 497 response_text = "Не смогла понять, какую задачу добавить." 498 elif parsed.action == "list_tasks": 499 if parsed.date_keyword is None: 500 tasks = task_manager.get_all_tasks(user_id) 501 incomplete_count = len(tasks) 502 scope = "all" 503 else: 504 date = parsed.date_keyword 505 tasks = task_manager.get_tasks_by_date(user_id, date) 506 incomplete_count = task_manager.get_incomplete_tasks_count( 507 user_id, date 508 ) 509 scope = "date" 510 511 if incomplete_count == 0: 512 if scope == "all": 513 response_text = "Невыполненных задач нет." 514 else: 515 response_text = f"На {parsed.date_label or parsed.date_keyword} невыполненных задач нет." 516 else: 517 if scope == "all": 518 prefix = f"У вас {incomplete_count} невыполненных задач: " 519 else: 520 prefix = f"На {parsed.date_label or parsed.date_keyword} у вас {incomplete_count} невыполненных задач: " 521 response_text = prefix + task_manager.format_task_list_truncated( 522 tasks, show_completed=False, max_chars=1024 - len(prefix) 523 ) 524 tts_formatted = task_manager.format_task_list_truncated( 525 tasks, show_completed=False, max_chars=1024 526 ) 527 if "\nи еще" in tts_formatted or tts_formatted.startswith( 528 "Не показано" 529 ): 530 summary = task_manager.get_incomplete_tasks_summary(user_id) 531 response_tts = ( 532 f"У вас {summary['total']} невыполненных задач. " 533 f"На сегодня {summary['today']}, на завтра {summary['tomorrow']}." 534 ) 535 else: 536 response_tts = tts_formatted 537 elif parsed.action == "help": 538 response_text = ( 539 "Я помогу вам управлять задачами. Вот что я умею:\n" 540 "• 'Добавь задачу купить молоко' - добавить задачу без срока\n" 541 "• 'Добавь задачу на завтра сходить к врачу' - добавить на завтра\n" 542 "• 'Добавь задачу на послезавтра ...' - добавить на послезавтра\n" 543 "• 'Какие задачи на сегодня?' - показать невыполненные задачи\n" 544 "• 'Дела на завтра' - показать задачи на завтра\n" 545 "• 'Что купить?' - найти задачи, начинающиеся с 'купить'\n" 546 "• 'Пометь задачу 3 выполненной' - отметить задачу как выполненную\n" 547 "• 'Удали задачу 2' / 'убери задачу 2' - удалить задачу\n" 548 "• 'Перенеси задачу 1 на завтра' - перенести одну задачу\n" 549 "• 'Перенеси задачи на завтра' - перенести все задачи" 550 ) 551 elif parsed.action == "what": 552 date = parsed.date_keyword 553 tasks = task_manager.get_tasks_by_date(user_id, date) 554 filtered_tasks = [ 555 t for t in tasks if t["text"].lower().startswith(parsed.keyword) 556 ] 557 incomplete_count = sum(1 for t in filtered_tasks if not t["completed"]) 558 559 if incomplete_count == 0: 560 response_text = f"На {parsed.date_label or parsed.date_keyword} нет задач, начинающихся с '{parsed.keyword}'." 561 else: 562 prefix = f"На {parsed.date_label or parsed.date_keyword} у вас {incomplete_count} задач: " 563 response_text = prefix + task_manager.format_task_list_truncated( 564 filtered_tasks, show_completed=False, max_chars=1024 - len(prefix) 565 ) 566 tts_formatted = task_manager.format_task_list_truncated( 567 filtered_tasks, show_completed=False, max_chars=1024 568 ) 569 if "\nи еще" in tts_formatted or tts_formatted.startswith( 570 "Не показано" 571 ): 572 summary = task_manager.get_incomplete_tasks_summary(user_id) 573 response_tts = ( 574 f"У вас {summary['total']} невыполненных задач. " 575 f"На сегодня {summary['today']}, на завтра {summary['tomorrow']}." 576 ) 577 else: 578 response_tts = tts_formatted 579 elif parsed.action == "complete": 580 if task_manager.complete_task(user_id, parsed.task_id): 581 response_text = f"Задача {parsed.task_id} отмечена как выполненная." 582 else: 583 response_text = f"Задача с номером {parsed.task_id} не найдена." 584 elif parsed.action == "delete": 585 if task_manager.delete_task(user_id, parsed.task_id): 586 response_text = f"Задача {parsed.task_id} удалена." 587 else: 588 response_text = f"Задача с номером {parsed.task_id} не найдена." 589 elif parsed.action == "no_task_id": 590 response_text = "Не смогла найти номер задачи." 591 elif parsed.action == "move_one": 592 date = parsed.date_keyword 593 if task_manager.move_task(user_id, parsed.task_id, date): 594 response_text = f"Задача {parsed.task_id} перенесена на {parsed.date_label or parsed.date_keyword}." 595 else: 596 response_text = f"Задача с номером {parsed.task_id} не найдена." 597 elif parsed.action == "move_all": 598 date = parsed.date_keyword 599 task_manager.move_tasks(user_id, date) 600 response_text = ( 601 f"Задачи перенесены на {parsed.date_label or parsed.date_keyword}." 602 ) 603 604 try: 605 get_langfuse().update_current_span( 606 metadata={ 607 "user_id": user_id, 608 "action": parsed.action if parsed else None, 609 "command": command, 610 "text_length": len(response_text), 611 "tts_length": len(response_tts) if response_tts else 0, 612 } 613 ) 614 except Exception: 615 pass 616 log.info(f"Отправляем ответ: {response_text}") 617 return build_response(request, response_text, tts=response_tts) 618 619 except Exception as e: 620 log.error(f"Ошибка обработки запроса: {e}") 621 return build_error_response()
Основной обработчик запросов от Яндекс Алисы.
Принимает POST-запрос от Диалогов, проверяет skill_id, парсит команду пользователя и выполняет соответствующее действие (добавление, просмотр, выполнение, удаление, перенос задач или вывод справки). Возвращает ответ в формате Яндекс Диалогов.
Args: request: Валидированный запрос от Яндекс Алисы
Returns: dict: Ответ навыка в формате Яндекс Диалогов
Raises: HTTPException: 403 если skill_id не совпадает
624@app.get("/") 625async def root(): 626 """ 627 Корневой эндпоинт для проверки доступности навыка. 628 629 Returns: 630 dict: Статус приложения и путь к файлу данных 631 """ 632 return { 633 "message": "Alice Todo Skill is running", 634 "status": "ok", 635 "data_file": DB_FILE, 636 }
Корневой эндпоинт для проверки доступности навыка.
Returns: dict: Статус приложения и путь к файлу данных
639@app.get("/health") 640async def health_check(): 641 """ 642 Эндпоинт проверки здоровья приложения. 643 644 Returns: 645 dict: Статус health, временная метка и количество пользователей 646 """ 647 return { 648 "status": "healthy", 649 "timestamp": datetime.now().isoformat(), 650 "users_count": task_manager.get_user_count(), 651 }
Эндпоинт проверки здоровья приложения.
Returns: dict: Статус health, временная метка и количество пользователей