import logging from logging.handlers import RotatingFileHandler import os import json import asyncio import httpx import re import smtplib import urllib3 import importlib.util from email.message import EmailMessage from ldap3 import Server, Connection, ALL from trueconf import Router, Message # --- SD WORKFLOW LOGGING SETUP --- log_dir = os.path.join(os.path.dirname(__file__), 'logs') os.makedirs(log_dir, exist_ok=True) sd_workflow_logger = logging.getLogger('sd_workflow') sd_workflow_logger.setLevel(logging.INFO) if not sd_workflow_logger.handlers: handler = RotatingFileHandler(os.path.join(log_dir, 'sd_workflow.log'), maxBytes=10*1024*1024, backupCount=5) formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s') handler.setFormatter(formatter) sd_workflow_logger.addHandler(handler) # ---------------------------------- # Абсолютный импорт конфигурации spec = importlib.util.spec_from_file_location("custom_config", os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "config", "config.py")) custom_config = importlib.util.module_from_spec(spec) spec.loader.exec_module(custom_config) globals().update(vars(custom_config)) from utils import states from utils.menu import MENU_TEXT from utils.texts import ( EMOJI_DIGITS, SD_MAIN_MENU_TEXT, SD_UNKNOWN_CMD_TEXT, SD_TEXT_REQUIRED, SD_CREATING, sd_ticket_created, sd_ticket_create_error, sd_system_error, sd_uploading_file, sd_file_attached, sd_file_upload_error, ) from utils.texts import UNKNOWN_MAIN_CMD_TEXT # 🔌 Импортируем централизованную функцию сбора статистики из main from utils.stats_logger import log_menu_stats # Таймаут ожидания вложений после текста (секунды) ATTACHMENT_WAIT_TIMEOUT = 2.0 SD_DISPATCHER_SYSTEM_PROMPT = ( "Ты — старший ИИ-диспетчер Service Desk холдинга. Твоя задача — сформулировать краткую и понятную тему ИТ-заявки по-русски (от 3 до 7 слов) на основе текста пользователя.\n" "ПРАВИЛА:\n" "1. Назови проблему и оборудование/программу (например: 'Самопроизвольное выключение ПК', 'Неисправность принтера', 'Сбой авторизации в почте').\n" "2. Отвечать на английском языке ЗАПРЕЩЕНО.\n" "3. Тебе категорически запрещено решать проблему или писать мануалы по настройке.\n" "4. Не пиши префиксы 'Тема:', 'Заголовок:' в ответе.\n" "5. Если в тексте только приветствие или мат — выведи ровно одно слово: СПАМ.") # Глушим системные предупреждения urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) logger = logging.getLogger(__name__) router = Router() sd_sessions = {} async def get_sd_templates(): url = f"{SD_URL}/api/v3/request_templates" headers = {"authtoken": SD_TOKEN, "Accept": "application/vnd.manageengine.sdp.v3+json"} input_data = {"list_info": {"row_count": 100, "start_index": 1}} params = {"input_data": json.dumps(input_data)} try: async with httpx.AsyncClient(verify=False) as client: response = await client.get(url, headers=headers, params=params, timeout=30.0) if response.status_code != 200: sd_workflow_logger.error(f"Ошибка API шаблонов: HTTP {response.status_code} - {response.text}") return [] data = response.json() all_templates = data.get("request_templates", []) pre_filtered = [ t for t in all_templates if t.get("status") == "ACTIVE" and not t.get("inactive", False) and str(t.get("id")) != "2" ] sd_workflow_logger.info(f"Шаблонов в системе: {len(all_templates)} | Активных: {len(pre_filtered)}") async def check_template(t): try: t_resp = await client.get(url + "/" + str(t["id"]), headers=headers, timeout=15.0) if t_resp.status_code == 200: detailed = t_resp.json().get("request_template", {}) return t if detailed.get("show_to_requester") is True else None except Exception as e: sd_workflow_logger.warning(f"Ошибка деталей шаблона ID {t['id']}: {e}") return None tasks = [check_template(t) for t in pre_filtered] results = await asyncio.gather(*tasks) active_templates = [t for t in results if t is not None] sd_workflow_logger.info(f"Шаблоны с show_to_requester=True: {len(active_templates)}") for temp in active_templates: sd_workflow_logger.info(f" ID: {temp['id']} | Name: {temp['name']}") return active_templates except Exception as e: sd_workflow_logger.error(f"Ошибка загрузки шаблонов: {e}") return [] # ================================================ # БИЗНЕС-ЛОГИКА # ================================================ def get_ad_user_sync(login: str): try: server = Server(AD_SERVER, get_info=ALL) conn = Connection(server, user=AD_USER, password=AD_PASSWORD, auto_bind=True) conn.search(search_base=AD_BASE, search_filter=f"(sAMAccountName={login})", attributes=["displayName", "mail", "l", "userAccountControl"]) if conn.entries: user = conn.entries[0] uac = user.userAccountControl.value if 'userAccountControl' in user else 0 return { "name": user.displayName.value if 'displayName' in user else login, "mail": user.mail.value if 'mail' in user else None, "city": user.l.value if 'l' in user else "Кемерово", "is_disabled": bool(uac & 2) } except Exception as e: logger.error(f"Ошибка LDAP: {e}") return None async def generate_smart_subject(text: str) -> str: import re, traceback, importlib.util, httpx spec = importlib.util.spec_from_file_location("custom_config", os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "config", "config.py")) custom_config = importlib.util.module_from_spec(spec) spec.loader.exec_module(custom_config) OLLAMA_SUBJECT_URL = custom_config.OLLAMA_SUBJECT_URL OLLAMA_SUBJECT_MODEL = custom_config.OLLAMA_SUBJECT_MODEL def has_chinese(t): return bool(re.search(r'[\u4e00-\u9fff]', str(t))) if not text or len(text.split()) < 2: return "Заявка из КЛЕВЕР" user_msg = f"Сформулируй краткую тему для следующей ИТ-заявки:\n{text}" messages = [ {"role": "system", "content": SD_DISPATCHER_SYSTEM_PROMPT}, {"role": "user", "content": user_msg} ] try: payload = {"model": OLLAMA_SUBJECT_MODEL, "messages": messages, "temperature": 0.0} sd_workflow_logger.info(f"📝 [Subject Gen] Payload: {payload}") endpoint = OLLAMA_SUBJECT_URL.replace("/api/generate", "/v1/chat/completions") async with httpx.AsyncClient() as client: response = await client.post(endpoint, json=payload, timeout=120.0) if response.status_code == 200: data = response.json() if "choices" in data and len(data["choices"]) > 0: res = data["choices"][0]["message"]["content"].strip(" .'\"`") res = re.sub(r'^(тема|заголовок|subject):\s*', '', res, flags=re.IGNORECASE).strip() if res and not has_chinese(res) and not any(m in res.upper() for m in ["СПАМ", "НЕ ПОНЯЛ", "Я ИСКУССТВЕННЫЙ", "Я МОДЕЛЬ"]): return res[0].upper() + res[1:] if len(res) > 0 else "Заявка из КЛЕВЕР" else: sd_workflow_logger.info(f"📝 [Subject Gen Error] Validation Failed: Received '{res}'") else: sd_workflow_logger.info(f"📝 [Subject Gen Error] Empty choices in LLM response") else: sd_workflow_logger.info(f"📝 [Subject Gen Error] API Error: Status {response.status_code} - {response.text}") except Exception as e: sd_workflow_logger.error(f"📝 [Subject Gen Error] Exception: {e}\n{traceback.format_exc()}") return "Заявка из КЛЕВЕР" async def create_ticket_in_sd(requester_email: str, subject: str, description: str, city: str, template_id: str = None): headers = {"authtoken": SD_TOKEN, "Accept": "application/vnd.manageengine.sdp.v3+json"} endpoint = f"{SD_URL}/api/v3/requests" html_desc = f"

{description.replace(chr(10), '
')}



Создано через TrueConf КЛЕВЕР

" async def _send_req(email, current_city): req_payload = {"request": {"subject": subject, "description": html_desc, "requester": {"email_id": email}, "udf_fields": {"udf_pick_301": current_city}}} if template_id: req_payload["request"]["request_template"] = {"id": template_id} async with httpx.AsyncClient(verify=False) as client: return await client.post(endpoint, headers=headers, data={"input_data": json.dumps(req_payload)}, timeout=15.0) resp = await _send_req(requester_email, city) data = resp.json() if resp.status_code in [200, 201] else {} if data.get("response_status", {}).get("status_code") != 2000: resp = await _send_req(requester_email, "Кемерово") data = resp.json() if resp.status_code in [200, 201] else {} if data.get("response_status", {}).get("status_code") == 2000: return data.get("request", {}).get("id") return None async def process_and_upload_file(file_id: str, filename: str, ticket_id: str): import io downloaded_path = None try: if isinstance(filename, bytes) or (isinstance(filename, str) and filename.startswith("b'")): safe_filename = f"attachment_{ticket_id}.dat" else: safe_filename = str(filename) downloaded_file = await states.tc_bot.download_file_by_id(file_id) if isinstance(downloaded_file, bytes): file_obj = io.BytesIO(downloaded_file) downloaded_path = None else: downloaded_path = str(downloaded_file) file_obj = open(downloaded_path, 'rb') try: upload_url = f"{SD_URL}/api/v3/requests/{ticket_id}/upload" async with httpx.AsyncClient(verify=False) as client: attach_resp = await client.put( upload_url, headers={"authtoken": SD_TOKEN}, files={'input_file': (safe_filename, file_obj, 'application/octet-stream')}, timeout=30.0 ) return attach_resp.status_code in [200, 201] finally: if not isinstance(downloaded_file, bytes): file_obj.close() except Exception as e: logger.error(f"Ошибка загрузки файла {filename}: {e}") return False finally: if downloaded_path and os.path.exists(downloaded_path): os.remove(downloaded_path) def _extract_text_from_content(content): """Извлечение текста из msg.content.""" if isinstance(content, dict): return content.get("text", "").strip() elif hasattr(content, "text"): return str(content.text).strip() return "" def _extract_attachments_from_content(content, is_attachment_type=False): """Извлечение file_id и file_name из msg.content.""" attachments = [] if isinstance(content, dict): atts = content.get("attachments") or content.get("files") single_fid = content.get("file_id") single_fname = content.get("file_name") if atts: if isinstance(atts, list): for a in atts: if isinstance(a, dict): fid = a.get("file_id") or a.get("id") fname = a.get("file_name") or a.get("name", "attachment") if fid: attachments.append({"file_id": fid, "file_name": fname}) elif isinstance(a, str): attachments.append({"file_id": a, "file_name": "attachment"}) elif isinstance(atts, dict): fid = atts.get("file_id") or atts.get("id") fname = atts.get("file_name") or atts.get("name", "attachment") if fid: attachments.append({"file_id": fid, "file_name": fname}) if single_fid: attachments.append({"file_id": single_fid, "file_name": single_fname or "attachment"}) elif is_attachment_type: if hasattr(content, "file_id"): fid = content.file_id fname = getattr(content, "file_name", "attachment") attachments.append({"file_id": fid, "file_name": fname}) return attachments async def _create_ticket_and_attach_files(user_id, msg_text, session, msg, login): """Логика генерации темы, отправки в SD и загрузки всех очередей вложений.""" session["step"] = "creating_ticket" template_id = session.get("template_id") sd_workflow_logger.info(f"[Create] Template ID: {template_id}") try: ad_user = await asyncio.to_thread(get_ad_user_sync, login) sd_workflow_logger.info(f"👤 [AD Lookup] User: {login} -> Found: {ad_user is not None}") sender_email = ad_user.get("mail") if ad_user else DEFAULT_REQUESTER city = ad_user.get("city") if ad_user else "Кемерово" subject = await generate_smart_subject(msg_text) sd_workflow_logger.info(f"📝 [Subject Gen] Text: {msg_text[:50]}... -> Subject: {subject}") ticket_id = await create_ticket_in_sd(sender_email, subject, msg_text, city, template_id) sd_workflow_logger.info(f"🎫 [Ticket Created] ID: {ticket_id}") if ticket_id: session["ticket_id"] = ticket_id # 1. Загрузка основных вложений (накопленных до и во время таймаута) queued_files = list(session.get("files_queue", [])) if queued_files: sd_workflow_logger.info(f"📎 [Upload] Uploading {len(queued_files)} file(s) to ticket #{ticket_id}") for f in queued_files: await process_and_upload_file(f['file_id'], f['file_name'], ticket_id) sd_workflow_logger.info(f"✅ [Upload] All queued file(s) uploaded") # 2. Загрузка вложений, прилетевших во время выполнения API-запросов (из "черной дыры") post_files = list(session.get("post_create_queue", [])) if post_files: sd_workflow_logger.info(f"📎 [Upload Extra] Uploading {len(post_files)} late file(s) to ticket #{ticket_id}") for f in post_files: await process_and_upload_file(f['file_id'], f['file_name'], ticket_id) session["files_queue"] = [] session["post_create_queue"] = [] session["step"] = "ticket_created" await msg.answer( sd_ticket_created(ticket_id, subject, msg_text), parse_mode="html" ) log_menu_stats(user_id, "Service Desk", f"Создание заявки #{ticket_id}") else: session["step"] = "need_text" await msg.answer(sd_ticket_create_error(), parse_mode="html") except Exception as e: logger.exception(f"Ошибка: {e}") sd_workflow_logger.error(f"❌ [SD Error] {e}") session["step"] = "need_text" await msg.answer(sd_system_error(), parse_mode="html") async def _wait_for_attachments_and_create(user_id, session): """Фоновый таймер с возможностью отмены (Debounce).""" try: await asyncio.sleep(ATTACHMENT_WAIT_TIMEOUT) except asyncio.CancelledError: return # Таймер сброшен, прерываем выполнение текущей таски if session.get("step") != "waiting_for_attachments": return sd_workflow_logger.info(f"⏰ [SD] Timeout reached. Initiating ticket creation for user {user_id}") msg = session.get("msg") login = session.get("login") if msg: await _create_ticket_and_attach_files(user_id, session["msg_text"], session, msg, login) @router.message() async def sd_module_handler(msg: Message): user_id = msg.from_user.id sd_workflow_logger.info(f"🚀 [SD Workflow Start] User: {user_id}") if states.get_state(user_id) != "SD_MODE": return # Инициализация шаблонов при первом входе в SD_MODE if user_id not in sd_sessions: sd_sessions[user_id] = {"step": "need_text", "files_queue": [], "post_create_queue": [], "templates": [], "template_names": {}, "template_menu_text": ""} templates = await get_sd_templates() if templates: sd_sessions[user_id]["templates"] = templates sd_sessions[user_id]["template_names"] = {str(i+1): str(t["id"]) for i, t in enumerate(templates)} template_menu_text = "" for i, t in enumerate(templates): template_menu_text += f"{EMOJI_DIGITS[str(i+1)]} - {t['name']}\n" sd_sessions[user_id]["template_menu_text"] = template_menu_text.strip() sd_workflow_logger.info(f"[Template] Loaded {len(templates)} templates for user {user_id}") sd_sessions[user_id]["step"] = "choose_template" await msg.answer( f"Выберите шаблон для заявки:\n\n" f"{template_menu_text}\n\n" f"Или отправьте пропустить для создания без шаблона.", parse_mode="html" ) return msg_text = _extract_text_from_content(msg.content) is_attachment = hasattr(msg.type, "name") and msg.type.name == "ATTACHMENT" inline_attachments = _extract_attachments_from_content(msg.content, is_attachment) sd_workflow_logger.info(f"🚀 [SD Debug] msg.type={msg.type}, content_type={type(msg.content).__name__}, msg_text='{msg_text}', is_attachment={is_attachment}, inline_attachments={len(inline_attachments)}") msg.handled = True login = user_id.split("@")[0] if "@" in user_id else user_id cmd = msg_text.lower() # 🔄 НОРМАЛИЗАЦИЯ КНОПОК ВК-ЭМОДЗИ for raw_num, emoji_num in EMOJI_DIGITS.items(): if cmd == emoji_num: cmd = raw_num break # --- ВЫХОД (Кнопки 0 и 9 возвращают в общее меню) --- if cmd in ["0", "9", "/start", "меню"]: log_menu_stats(user_id, "Service Desk", "Выход в главное меню") states.clear_state(user_id) sd_sessions.pop(user_id, None) await msg.answer(MENU_TEXT, parse_mode="html") return if user_id not in sd_sessions: sd_sessions[user_id] = {"step": "need_text", "files_queue": [], "post_create_queue": [], "templates": [], "template_names": {}} session = sd_sessions[user_id] # --- ШАГ 0: ВЫБОР ШАБЛОНА --- if session["step"] == "choose_template": if cmd in ["skip", "пропустить", "без шаблона"]: sd_workflow_logger.info("[Template] User skipped template selection") session["step"] = "need_text" await msg.answer(SD_MAIN_MENU_TEXT, parse_mode="html") return if cmd in session.get("template_names", {}): template_id = session["template_names"][cmd] template_name = next((t["name"] for t in session.get("templates", []) if str(t["id"]) == str(template_id)), "Неизвестно") sd_workflow_logger.info(f"[Template] User selected: {template_id} ({template_name})") session["template_id"] = template_id session["step"] = "need_text" await msg.answer( f"Выбран шаблон: {template_name}\n\n" f"Теперь опишите проблему текстом (можно прикрепить файлы):", parse_mode="html" ) return else: await msg.answer( f"Неверный номер. Выберите из списка:\n\n" f"{session.get('template_menu_text', '')}\n\n" f"Или отправьте пропустить для создания без шаблона.", parse_mode="html" ) return # --- ШАГ 1: ОЖИДАНИЕ ТЕКСТА --- if session["step"] == "need_text": if not msg_text.strip() and not inline_attachments: await msg.answer(SD_UNKNOWN_CMD_TEXT, parse_mode="html") return # Если только вложение без текста — сохраняем и ждем текст if not msg_text.strip() and inline_attachments: session["files_queue"].extend(inline_attachments) sd_workflow_logger.info(f"📎 [SD] Saved {len(inline_attachments)} attachment(s), waiting for text") await msg.answer(SD_TEXT_REQUIRED, parse_mode="html") return # Если пришел текст (с вложениями или без) if inline_attachments: session["files_queue"].extend(inline_attachments) sd_workflow_logger.info(f"📎 [SD] Added {len(inline_attachments)} inline attachment(s) to queue") # Настраиваем параметры сессии для ожидания session["step"] = "waiting_for_attachments" session["msg_text"] = msg_text session["msg"] = msg session["user_id"] = user_id session["login"] = login # Сразу выдаем ОДНО сообщение пользователю, чтобы он видел реакцию бота await msg.answer(SD_CREATING, parse_mode="html") # Запускаем фоновый таймер с механизмом сброса (Debounce) if "timer_task" in session and not session["timer_task"].done(): session["timer_task"].cancel() session["timer_task"] = asyncio.create_task(_wait_for_attachments_and_create(user_id, session)) sd_workflow_logger.info(f"⏳ [SD] Timer started: waiting {ATTACHMENT_WAIT_TIMEOUT}s for potentially more attachments...") return # --- ШАГ 2: РЕЖИМ ОЖИДАНИЯ ДОП. ВЛОЖЕНИЙ (Сброс таймаута) --- if session["step"] == "waiting_for_attachments": if inline_attachments: session["files_queue"].extend(inline_attachments) sd_workflow_logger.info(f"📎 [SD] Received {len(inline_attachments)} more files. Total queue: {len(session['files_queue'])}") if msg_text.strip(): session["msg_text"] = session["msg_text"] + " " + msg_text.strip() sd_workflow_logger.info(f"📝 [SD] Additional text appended") # Перезапускаем таймер (пользователь активен, сдвигаем окно создания вперед) if "timer_task" in session and not session["timer_task"].done(): session["timer_task"].cancel() session["timer_task"] = asyncio.create_task(_wait_for_attachments_and_create(user_id, session)) sd_workflow_logger.info(f"🔄 [SD] Timer reset due to user activity. Waiting another {ATTACHMENT_WAIT_TIMEOUT}s") return # --- ШАГ 3: ЗАЯВКА СОЗДАЕТСЯ (Перехват файлов из "черной дыры") --- if session["step"] == "creating_ticket": if inline_attachments: if "post_create_queue" not in session: session["post_create_queue"] = [] session["post_create_queue"].extend(inline_attachments) sd_workflow_logger.info(f"📥 [SD Hole-Fix] Captured {len(inline_attachments)} file(s) DURING ticket creation API call.") return # --- ШАГ 4: ПРИЕМ ДОП. ФАЙЛОВ К УЖЕ СОЗДАННОЙ ЗАЯВКЕ --- if session["step"] == "ticket_created": if inline_attachments: ticket_id = session["ticket_id"] for att in inline_attachments: await msg.answer(sd_uploading_file(ticket_id), parse_mode="html") ok = await process_and_upload_file(att['file_id'], att['file_name'], ticket_id) sd_workflow_logger.info(f"📎 [Late File Upload] Name: {att['file_name']} -> Success: {ok}") if ok: log_menu_stats(user_id, "Service Desk", f"Добавление файла к заявке #{ticket_id}") await msg.answer(sd_file_attached(), parse_mode="html") else: await msg.answer(sd_file_upload_error(), parse_mode="html") elif msg_text.strip(): await msg.answer(UNKNOWN_MAIN_CMD_TEXT, parse_mode="html") return