#!/usr/bin/env python3 """ email_classifier.py — классификация писем локальной LLM (Qwen3:8b через Ollama). Читает email.md файлы архива, для каждого письма БЕЗ поля `classification` вызывает Qwen3:8b (localhost:11434), получает JSON с тегом и обоснованием, записывает в frontmatter: classification: info|urgent|task|meeting|task,meeting|unclassified classification_reason: "краткое обоснование на русском" meeting_datetime: "YYYY-MM-DD HH:MM" (только для meeting) Трекинг обработанных — по наличию `classification` в frontmatter (D2): повторный запуск пропускает уже обработанные письма. Запуск: python3 scripts/email_classifier.py # новые письма (свежие первыми) python3 scripts/email_classifier.py --limit 10 python3 scripts/email_classifier.py --force # переклассифицировать всё python3 scripts/email_classifier.py --folder INBOX """ import argparse import json import re import sys import time import urllib.error import urllib.request from datetime import datetime from pathlib import Path # Конфигурация EMAIL_ROOT = Path("/opt/hermes/email") OLLAMA_URL = "http://localhost:11434/api/generate" OLLAMA_MODEL = "qwen3:8b-nothink" # текстовая задача (без think-токенов — быстрее) LLM_TIMEOUT = 60 # секунд на один запрос (классификация длиннее контактов) MAX_BODY_CHARS = 5000 # как в contacts_extractor TEMP = 0.1 # Теги классификации (валидные значения поля classification) VALID_TAGS = {"info", "urgent", "task", "meeting"} PROMPT_TEMPLATE = """Ты — классификатор входящей почты. Определи тип письма по его тексту. Возможные типы (можно комбинировать через запятую): - info: информационное письмо, не требует действий (новости, рассылки, отчёты для сведения) - urgent: требует срочного ответа/действия сегодня (горящие сроки, просьбы ответить) - task: содержит поручение/задачу, которую нужно выполнить (что-то сделать, подготовить, прислать) - meeting: содержит приглашение на встречу/совещание/созвон, или просьбу назначить встречу Правила: - Если письмо содержит и задачу, и встречу — верни "task,meeting" - Если явно не указано — лучше info, чем ложное срабатывание - Для meeting попробуй извлечь дату и время из текста (формат "YYYY-MM-DD HH:MM", время в 24-часовом формате, например "2026-09-15 11:00"). Если дата не указана — null. Верни ТОЛЬКО JSON, без пояснений: {{"classification": "info", "reason": "1-2 предложения на русском, почему такой тег", "meeting_datetime": null}} Тема письма: {subject} Отправитель: {sender} Текст письма: {body} """ def parse_email_md(path): """Прочитать email.md, вернуть (headers_dict, body_text, raw_content, fm_end).""" content = path.read_text(encoding="utf-8", errors="replace") match = re.match(r"^---\s*\n(.*?)\n---\s*\n(.*)", content, re.DOTALL) if match: yaml_block = match.group(1) body = match.group(2).strip() fm_end = match.end(1) # позиция конца YAML-блока (перед закрывающим ---) headers = {} for line in yaml_block.split("\n"): m = re.match(r"^(\w[\w_-]*)\s*:\s*(.*)$", line) if m: headers[m.group(1)] = m.group(2).strip() else: headers, body, fm_end = {}, content.strip(), None return headers, body, content, fm_end def clean_body(body): """Очистка тела письма — переиспользуем логику contacts_extractor.""" # Удаляем <#part ...> блоки и HTML-теги body = re.sub(r"<#part[^>]*>", "", body) body = re.sub(r"<#/part>", "", body) body = re.sub(r"<[^>]+>", "", body) body = re.sub(r"\(mailto:[^)]+\)", "", body) # Unicode-пробелы → обычные body = re.sub(r"[\u00a0\u2000-\u200f\u2028-\u202f\u2060]+", " ", body) # Трекинг-ссылки body = re.sub(r"https?://tn-eoc\.[^\s]+", "", body) body = re.sub(r"https?://[^\s]+\?utm_[^\s]+", "", body) # Цитируемая переписка — отрезаем от самого раннего маркера quote_patterns = [ r"^[\s]*_{4,}\s*$", r"От:.*\n[\s]*Отправлено:", r"^[\s]*From:.*\n[\s]*Sent:", r"—+.*Forwarded.*—+", r"—+.*Пересылаемое.*—+", r"—+.*Original Message.*—+", r">.*\bwrote:", ] earliest_pos = len(body) for qp in quote_patterns: for m in re.finditer(qp, body, re.MULTILINE): if m.start() < earliest_pos: earliest_pos = m.start() if earliest_pos < len(body): body = body[:earliest_pos].strip() else: tail = body[-500:] if len(body) > 500 else body for pattern in [r"От:", r"Отправлено:", r"From:", r"Sent:", r"Кому:", r"To:", r"Тема:", r"Subject:"]: m2 = re.search(pattern, tail) if m2: offset = len(body) - len(tail) + m2.start() body = body[:offset].strip() break lines = [l for l in body.split("\n") if not re.match(r"^\s*>", l)] body = re.sub(r"\n{3,}", "\n\n", "\n".join(lines)) return body.strip() def yaml_quote(v): """YAML-значение: обернуть в двойные кавычки при спецсимволах.""" v = str(v) if v == "": return '""' if re.search(r'[:#\[\]{}&*!|>\'"%@`\n]|^\s|\s$', v): return '"' + v.replace("\\", "\\\\").replace('"', '\\"') + '"' return v def call_llm(subject, sender, body_text, max_retries=2): """Вызвать Qwen через Ollama, вернуть dict или None.""" body_text = clean_body(body_text)[:MAX_BODY_CHARS] prompt = PROMPT_TEMPLATE.format(subject=subject or "(без темы)", sender=sender or "?", body=body_text) for attempt in range(max_retries + 1): if attempt > 0: time.sleep(1) payload = json.dumps({ "model": OLLAMA_MODEL, "prompt": prompt, "stream": False, "options": {"temperature": TEMP, "num_predict": 512}, }).encode("utf-8") req = urllib.request.Request(OLLAMA_URL, data=payload, headers={"Content-Type": "application/json"}, method="POST") try: resp = urllib.request.urlopen(req, timeout=LLM_TIMEOUT) data = json.loads(resp.read().decode("utf-8")) response_text = data.get("response", "").strip() except (urllib.error.URLError, json.JSONDecodeError, TimeoutError) as e: if attempt < max_retries: continue print(f" ⚠ LLM error: {e}", file=sys.stderr) return None if not response_text: if attempt < max_retries: continue print(f" ⚠ LLM empty response", file=sys.stderr) return None # Парсим: весь ответ как JSON или { ... } внутри try: return json.loads(response_text) except json.JSONDecodeError: pass brace_depth, json_start = 0, None for i, ch in enumerate(response_text): if ch == "{": if brace_depth == 0: json_start = i brace_depth += 1 elif ch == "}": brace_depth -= 1 if brace_depth == 0 and json_start is not None: try: return json.loads(response_text[json_start:i + 1]) except json.JSONDecodeError: pass json_start = None if attempt < max_retries: continue print(f" ⚠ LLM JSON parse error: {response_text[:300]}", file=sys.stderr) return None def normalize_classification(raw): """Привести теги к валидному виду (через запятую), вернуть (tags_str, reason, meeting_dt).""" tags_raw = raw.get("classification") or raw.get("tags") or "" if isinstance(tags_raw, list): tags = [t.strip().lower() for t in tags_raw if isinstance(t, str)] else: tags = [t.strip().lower() for t in str(tags_raw).split(",") if t.strip()] # Оставляем только валидные теги tags = [t for t in tags if t in VALID_TAGS] if not tags: return "unclassified", (raw.get("reason") or "").strip(), None tags = sorted(set(tags)) # детерминированный порядок reason = (raw.get("reason") or raw.get("classification_reason") or "").strip() meeting_dt = None if "meeting" in tags: mdt = raw.get("meeting_datetime") if mdt: s = str(mdt).strip() m = re.match(r"^(\d{4}-\d{2}-\d{2})[T ](\d{1,2}:\d{2})", s) if m: meeting_dt = f"{m.group(1)} {m.group(2)}" return ",".join(tags), reason, meeting_dt def add_to_frontmatter(content, fm_end, fields): """Добавить поля YAML в frontmatter (перед закрывающим ---).""" add_lines = [] for k, v in fields: if v is None or v == "": continue add_lines.append(f"{k}: {yaml_quote(v)}") if not add_lines: return content before = content[:fm_end] after = content[fm_end:] return before + "\n" + "\n".join(add_lines) + after def find_email_md_files(root, folder=None): """Найти email.md, опционально в конкретной папке (префикс пути).""" files = [] for p in sorted(root.rglob("email.md")): if folder: rel = p.relative_to(root) if not rel.parts[0] == folder: continue files.append(p) return files def email_sort_key(path): """Свежие письма первыми: по (году, месяцу) из пути + UID (число).""" parts = path.parts # Ищем в частях пути год (4 цифры), месяц (2), UID (число-каталог > 100) year = next((int(p) for p in parts if re.fullmatch(r"\d{4}", p)), 0) month = next((int(p) for p in parts if re.fullmatch(r"\d{2}", p) and 1 <= int(p) <= 12), 0) uid = next((int(p) for p in parts if p.isdigit() and int(p) > 100), 0) # fallback: mtime файла if not year: try: return (-float(path.stat().st_mtime),) except OSError: return (0,) return (-year, -month, -uid) def main(): ap = argparse.ArgumentParser(description="Классификация писем через Qwen3:8b (Ollama)") ap.add_argument("--limit", type=int, default=0, help="Максимум писем за проход (0 = все)") ap.add_argument("--force", action="store_true", help="Переклассифицировать даже обработанные") ap.add_argument("--folder", default=None, help="Только письма из конкретной папки (INBOX)") args = ap.parse_args() print(f"Классификатор: {OLLAMA_MODEL} ({OLLAMA_URL})") files = find_email_md_files(EMAIL_ROOT, args.folder) print(f"Найдено email.md: {len(files)}") pending = [] for p in files: headers, body, content, fm_end = parse_email_md(p) if not args.force and headers.get("classification"): continue pending.append((p, headers, body, content, fm_end)) # Свежие первыми pending.sort(key=lambda x: email_sort_key(x[0])) print(f"Классифицировать: {len(pending)}") if args.limit > 0: pending = pending[:args.limit] processed = 0 for p, headers, body, content, fm_end in pending: try: subject = headers.get("subject", "") sender = headers.get("from", "") result = call_llm(subject, sender, body) if result is None: tags, reason, meeting_dt = "unclassified", "LLM не ответила", None else: tags, reason, meeting_dt = normalize_classification(result) fields = [] if headers.get("classification"): # --force: обновляем, но поля уже есть — перезапишем через добавление fields.append(("classification", tags)) fields.append(("classification_reason", reason)) if meeting_dt: fields.append(("meeting_datetime", meeting_dt)) else: fields.append(("classification", tags)) fields.append(("classification_reason", reason)) if meeting_dt: fields.append(("meeting_datetime", meeting_dt)) new_content = add_to_frontmatter(content, fm_end, fields) if new_content != content: p.write_text(new_content, encoding="utf-8") print(f" ✓ {p.parent.parent.parent.name}/{p.parent.name}/{headers.get('subject','')[:50]!r} → {tags}") processed += 1 except Exception as e: print(f" ✗ {p}: {e}", file=sys.stderr) print(f"\nГотово. Обработано: {processed}") if __name__ == "__main__": main()