Files
email-assistant/scripts/email_classifier.py
T
hermes 7550aff102 feat: классификация писем Qwen3:8b + обработчики (urgent→TG, task→VTODO, meeting→VEVENT)
- mail_archive.py: фикс вложений himalaya --dir → --downloads-dir; attachments/ только при has_attachment; идемпотентно
- email_classifier.py: Qwen3:8b (Ollama) → теги info/urgent/task/meeting в frontmatter email.md
- email_handlers.py: urgent→Telegram (Bot API+SOCKS5), task→Radicale VTODO, meeting→Radicale VEVENT; handled_* идемпотентность; --dry-run
- .gitignore: игнор *.env.bak*
- openspec: чейндж email-classification-handlers (в работе)
- STATUS/TODO/WALKTHROUGH: прогресс сессии, подводные камни Radicale (http.client, Depth:1, формат дат)
2026-09-13 20:35:58 +00:00

323 lines
14 KiB
Python

#!/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()