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, формат дат)
This commit is contained in:
2026-09-13 20:35:58 +00:00
parent 8ea022f5c0
commit 7550aff102
22 changed files with 1728 additions and 76 deletions
+323
View File
@@ -0,0 +1,323 @@
#!/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()
+456
View File
@@ -0,0 +1,456 @@
#!/usr/bin/env python3
"""
email_handlers.py — обработчики по тегам классификации писем.
Сканирует email.md архива, для писем с тегом classification и без соответствующего
поля handled_* в frontmatter выполняет обработчик:
urgent → Telegram (через Bot API + SOCKS5-туннель; from/subject/превью)
task → Radicale CalDAV: VTODO в календарь «Задачи» (SUMMARY=тема,
DESCRIPTION=ссылка на email.md, при наличии даты DTSTART/DUE)
meeting → Radicale CalDAV: VEVENT в календарь «Рабочий» (SUMMARY=тема,
DTSTART из meeting_datetime или ближайший рабочий день 11:00)
info → ничего (только тег в frontmatter)
После успешной обработки в frontmatter пишется handled_urgent/handled_task/
handled_meeting: true — повторный запуск не создаёт дубликатов (идемпотентность).
Секреты — только из .env (рядом со скриптом):
RADICALE_URL / RADICALE_USER / RADICALE_PASS — доступ к Radicale
VESTI_BOT_TOKEN (или TELEGRAM_BOT_TOKEN) — токен бота Telegram
TG_PROXY — SOCKS5 до Bot API (по умолчанию socks5://127.0.0.1:1080)
TELEGRAM_CHAT_ID — куда слать urgent (по умолчанию @dedinit_vesti)
Запуск:
python3 scripts/email_handlers.py # все необработанные
python3 scripts/email_handlers.py --limit 10
python3 scripts/email_handlers.py --folder INBOX
python3 scripts/email_handlers.py --dry-run # показать, что бы сделал
"""
import argparse
import http.client
import json
import os
import re
import sys
import urllib.error
import urllib.request
from datetime import datetime, timedelta
from functools import lru_cache
from pathlib import Path
try:
from dotenv import load_dotenv
except ImportError:
load_dotenv = None
# Каталог скрипта → .env рядом с проектом
BASE_DIR = Path(__file__).resolve().parents[1]
if load_dotenv:
load_dotenv(BASE_DIR / ".env", override=False)
EMAIL_ROOT = Path(os.getenv("EMAIL_ROOT", "/opt/hermes/email"))
# --- Radicale (CalDAV) ---
RADICALE_URL = os.getenv("RADICALE_URL", "http://127.0.0.1:5232").rstrip("/")
RADICALE_USER = os.getenv("RADICALE_USER", "estorozhenko")
RADICALE_PASS = os.getenv("RADICALE_PASS", "")
# Календари (percent-encoded, «Задачи» и «Рабочий» — кириллица)
TASKS_CAL = os.getenv("RADICALE_TASKS_CAL", "%D0%97%D0%B0%D0%B4%D0%B0%D1%87%D0%B8") # Задачи
WORK_CAL = os.getenv("RADICALE_WORK_CAL", "%D0%A0%D0%B0%D0%B1%D0%BE%D1%87%D0%B8%D0%B9") # Рабочий
# --- Telegram (Bot API через SOCKS5) ---
TG_TOKEN = os.getenv("VESTI_BOT_TOKEN") or os.getenv("TELEGRAM_BOT_TOKEN") or ""
TG_PROXY = os.getenv("TG_PROXY", "socks5://127.0.0.1:1080")
TG_CHAT_ID = os.getenv("TELEGRAM_CHAT_ID", "@dedinit_vesti")
TG_API = "https://api.telegram.org"
# --- Общие ---
LLM_TIMEOUT = 20
MAX_PREVIEW_CHARS = 400 # превью письма для Telegram
TG_TIMEOUT = 20
# ---------------------------------------------------------------------------
# Frontmatter
# ---------------------------------------------------------------------------
def parse_email_md(path):
"""Прочитать email.md, вернуть (headers, body, 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)
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 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 add_to_frontmatter(content, fm_end, fields):
"""Добавить поля в 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 mark_handled(path, tag):
"""Пометить письмо handled_<tag>: true. Возвращает True при изменении."""
headers, body, content, fm_end = parse_email_md(path)
if fm_end is None:
return False
key = f"handled_{tag}"
if headers.get(key) == "true":
return False
new_content = add_to_frontmatter(content, fm_end, [(key, "true")])
if new_content != content:
path.write_text(new_content, encoding="utf-8")
return True
# ---------------------------------------------------------------------------
# Radicale (CalDAV)
# ---------------------------------------------------------------------------
def _caldav(path: str, method="GET", body=None, content_type=None):
"""Базовый HTTP к Radicale с Basic-auth через http.client.
path — ПОЛНЫЙ URL (например http://127.0.0.1:5232/estorozhenko/...).
urllib.request не умеет URL с percent-encoded кириллицей в пути
(Errno -2 Name or service not known), поэтому используем http.client
напрямую — он корректно работает с encoded path.
Возвращает (status, text).
"""
import base64
from urllib.parse import urlsplit
if not path.startswith("http"):
path = RADICALE_URL + path
parsed = urlsplit(path)
conn = http.client.HTTPConnection(parsed.hostname, parsed.port, timeout=LLM_TIMEOUT)
full_path = parsed.path + (("?" + parsed.query) if parsed.query else "")
headers = {}
if content_type:
headers["Content-Type"] = content_type
if method == "PROPFIND":
headers["Depth"] = "1"
cred = base64.b64encode(f"{RADICALE_USER}:{RADICALE_PASS}".encode()).decode()
headers["Authorization"] = f"Basic {cred}"
try:
conn.request(method, full_path, body=body, headers=headers)
resp = conn.getresponse()
return resp.status, resp.read().decode("utf-8", errors="replace")
except Exception as e:
return 0, str(e)
finally:
conn.close()
def ics_escape(s):
"""Экранирование значений iCalendar (RFC 5545): backslash, semicolon, comma, переносы."""
s = str(s).replace("\\", "\\\\").replace(";", "\\;").replace(",", "\\,")
return s.replace("\r\n", "\\n").replace("\n", "\\n")
@lru_cache(maxsize=1)
def find_calendar_url():
"""Определить URL коллекций «Задачи» и «Рабочий» из PROPFIND (xml.etree)."""
status, text = _caldav(f"/{RADICALE_USER}/", "PROPFIND", body=b"", content_type="application/xml")
found: dict = {"Задачи": None, "Рабочий": None}
if status != 207:
return found
try:
import xml.etree.ElementTree as ET
root = ET.fromstring(text)
# Пространства имён: DAV: (по умолчанию), C: — caldav
ns = {"d": "DAV:", "c": "urn:ietf:params:xml:ns:caldav"}
for resp in root.findall("d:response", ns):
href_el = resp.find("d:href", ns)
if href_el is None:
continue
href = (href_el.text or "").strip()
rt = resp.find(".//d:resourcetype", ns)
if rt is None:
continue
is_cal = rt.find("c:calendar", ns) is not None
if not is_cal:
continue
import urllib.parse
dec = urllib.parse.unquote(href)
for name, key in [("Задачи", "Задачи"), ("Рабочий", "Рабочий")]:
if key in dec and found[name] is None:
found[name] = RADICALE_URL + href
except Exception as e:
print(f" ⚠ find_calendar_url: {e}", file=sys.stderr)
return found
def vtodo(uid, summary, description, due_dt=None):
"""Сформировать VTODO (iCalendar). due_dt: 'YYYY-MM-DD HH:MM' или None."""
now = datetime.now().strftime("%Y%m%dT%H%M%S")
lines = [
"BEGIN:VCALENDAR",
"VERSION:2.0",
"PRODID:-//email-assistant//VTODO//RU",
"BEGIN:VTODO",
f"UID:{uid}@email-assistant",
f"DTSTAMP:{now}",
f"SUMMARY:{ics_escape(summary)}",
f"DESCRIPTION:{ics_escape(description)}",
]
if due_dt:
lines.append(f"DUE:{due_dt.replace(' ', 'T')}:00")
lines += ["STATUS:NEEDS-ACTION", "END:VTODO", "END:VCALENDAR"]
return "\r\n".join(lines)
def vevent(uid, summary, description, start_dt, duration_min=60):
"""Сформировать VEVENT. start_dt: 'YYYY-MM-DD HH:MM'."""
now = datetime.now().strftime("%Y%m%dT%H%M%S")
start = datetime.strptime(start_dt, "%Y-%m-%d %H:%M")
start_ics = start.strftime("%Y%m%dT%H%M%S")
end_ics = (start + timedelta(minutes=duration_min)).strftime("%Y%m%dT%H%M%S")
lines = [
"BEGIN:VCALENDAR",
"VERSION:2.0",
"PRODID:-//email-assistant//VEVENT//RU",
"BEGIN:VEVENT",
f"UID:{uid}@email-assistant",
f"DTSTAMP:{now}",
f"SUMMARY:{ics_escape(summary)}",
f"DESCRIPTION:{ics_escape(description)}",
f"DTSTART:{start_ics}",
f"DTEND:{end_ics}",
"END:VEVENT",
"END:VCALENDAR",
]
return "\r\n".join(lines)
def next_workday_1100(now=None):
"""Ближайший будний день (пн-пт) в 11:00. now: datetime."""
now = now or datetime.now()
d = now
while d.weekday() >= 5: # сб/вс
d += timedelta(days=1)
return d.strftime("%Y-%m-%d") + " 11:00"
def handle_task(path, headers, body):
"""Создать VTODO в Radicale «Задачи». Возвращает (ok, detail)."""
summary = (headers.get("subject") or "(без темы)").strip()
# UID стабильный: от пути письма
rel = path.relative_to(EMAIL_ROOT) if EMAIL_ROOT in path.parents else path
uid = re.sub(r"[^a-zA-Z0-9]+", "-", str(rel)).strip("-")
description = f"Из письма: {path}"
due = None
mdt = headers.get("meeting_datetime") or headers.get("date")
if mdt:
mdt = mdt.replace("T", " ")[:16]
if re.match(r"^\d{4}-\d{2}-\d{2} \d{2}:\d{2}$", mdt):
due = mdt
# URL коллекции «Задачи»
cals = find_calendar_url()
cal_url = cals.get("Задачи")
if not cal_url:
return False, "Коллекция «Задачи» не найдена в PROPFIND"
ics = vtodo(uid, summary, description, due)
resp_status, resp_text = _caldav(cal_url + uid + ".ics", "PUT", body=ics.encode("utf-8"),
content_type="text/calendar; charset=utf-8")
if resp_status in (200, 201, 204):
return True, f"VTODO создан ({resp_status})"
return False, f"PUT {resp_status}: {resp_text[:200]}"
def handle_meeting(path, headers, body):
"""Создать VEVENT в Radicale «Рабочий». Возвращает (ok, detail)."""
summary = (headers.get("subject") or "(без темы)").strip()
rel = path.relative_to(EMAIL_ROOT) if EMAIL_ROOT in path.parents else path
uid = re.sub(r"[^a-zA-Z0-9]+", "-", str(rel)).strip("-")
description = f"Из письма: {path}"
start = None
mdt = headers.get("meeting_datetime")
if mdt:
mdt = mdt.replace("T", " ")[:16]
if re.match(r"^\d{4}-\d{2}-\d{2} \d{2}:\d{2}$", mdt):
start = mdt
if not start:
start = next_workday_1100()
cals = find_calendar_url()
cal_url = cals.get("Рабочий")
if not cal_url:
return False, "Коллекция «Рабочий» не найдена в PROPFIND"
ics = vevent(uid, summary, description, start)
resp_status, resp_text = _caldav(cal_url + uid + ".ics", "PUT", body=ics.encode("utf-8"),
content_type="text/calendar; charset=utf-8")
if resp_status in (200, 201, 204):
return True, f"VEVENT создан ({resp_status})"
return False, f"PUT {resp_status}: {resp_text[:200]}"
# ---------------------------------------------------------------------------
# Telegram (Bot API через SOCKS5)
# ---------------------------------------------------------------------------
def tg_client():
"""httpx.Client с SOCKS5-прокси. Путь прокси из TG_PROXY."""
try:
import httpx
except ImportError:
return None
return httpx.Client(proxy=TG_PROXY, timeout=TG_TIMEOUT)
def tg_call(method, **params):
"""Вызвать метод Bot API. Возвращает result dict. Ошибки -> исключение."""
if not TG_TOKEN:
raise RuntimeError("Нет токена Telegram (VESTI_BOT_TOKEN/TELEGRAM_BOT_TOKEN не задан)")
client = tg_client()
if client is None:
raise RuntimeError("httpx не установлен — нужен для Telegram")
try:
with client:
r = client.post(f"{TG_API}/bot{TG_TOKEN}/{method}", json=params, timeout=TG_TIMEOUT)
except Exception as e:
raise RuntimeError(f"Сеть/прокси до Bot API: {e}") from e
if r.status_code != 200:
try:
desc = r.json().get("description", "")
except Exception:
desc = r.text[:200]
raise RuntimeError(f"Bot API {method}: HTTP {r.status_code} {desc}")
data = r.json()
if not data.get("ok"):
raise RuntimeError(f"Bot API {method}: {data.get('description','')}")
return data.get("result", {})
def handle_urgent(path, headers, body):
"""Отправить уведомление в Telegram. Возвращает (ok, detail)."""
subject = (headers.get("subject") or "(без темы)").strip()
sender = (headers.get("from") or "?").strip()
preview = body.strip()
if len(preview) > MAX_PREVIEW_CHARS:
preview = preview[:MAX_PREVIEW_CHARS].rstrip() + "…"
text = (
f"⚠️ СРОЧНОЕ письмо\n\n"
f"От: {sender}\n"
f"Тема: {subject}\n\n"
f"{preview}\n\n"
f"Письмо: file://{path}"
)
try:
result = tg_call("sendMessage", chat_id=TG_CHAT_ID, text=text,
link_preview_options={"is_disabled": True})
msg_id = result.get("message_id")
return True, f"Отправлено в {TG_CHAT_ID} (msg_id={msg_id})"
except Exception as e:
return False, f"Telegram: {e}"
# ---------------------------------------------------------------------------
# Основной проход
# ---------------------------------------------------------------------------
def find_email_md_files(root, folder=None):
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):
"""Свежие письма первыми."""
import re as _re
parts = path.parts
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)
if not year:
try:
return (-float(path.stat().st_mtime),)
except OSError:
return (0,)
return (-year, -month, -uid)
HANDLERS = {
"urgent": handle_urgent,
"task": handle_task,
"meeting": handle_meeting,
}
def main():
ap = argparse.ArgumentParser(description="Обработчики по тегам классификации писем")
ap.add_argument("--limit", type=int, default=0, help="Максимум писем за проход")
ap.add_argument("--folder", default=None, help="Только папка (INBOX)")
ap.add_argument("--dry-run", action="store_true", help="Не писать, только показать")
args = ap.parse_args()
files = find_email_md_files(EMAIL_ROOT, args.folder)
pending = []
for p in files:
headers, body, content, fm_end = parse_email_md(p)
cls = headers.get("classification", "")
if not cls:
continue
tags = [t.strip() for t in cls.split(",") if t.strip()]
need = [t for t in tags if t in HANDLERS and headers.get(f"handled_{t}") != "true"]
if need:
pending.append((p, headers, body, need))
pending.sort(key=lambda x: email_sort_key(x[0]))
print(f"Обработать: {len(pending)}")
if args.limit > 0:
pending = pending[:args.limit]
results = {"urgent": 0, "task": 0, "meeting": 0, "errors": 0}
for p, headers, body, need in pending:
for tag in need:
handler = HANDLERS[tag]
try:
ok, detail = handler(p, headers, body)
if ok:
if not args.dry_run:
mark_handled(p, tag)
results[tag] += 1
print(f" ✓ [{tag}] {p.parent.parent.parent.name}/{p.parent.name}: {detail}")
else:
results["errors"] += 1
print(f" ✗ [{tag}] {p.parent.parent.parent.name}/{p.parent.name}: {detail}", file=sys.stderr)
except Exception as e:
results["errors"] += 1
print(f" ✗ [{tag}] {p}: {e}", file=sys.stderr)
print(f"\nГотово: urgent={results['urgent']}, task={results['task']}, meeting={results['meeting']}, ошибок={results['errors']}")
if args.dry_run:
print("(dry-run: ничего не записано и не отправлено)")
if __name__ == "__main__":
main()
+30 -8
View File
@@ -398,18 +398,34 @@ def make_email_md(meta, extra_headers, body, folder_name):
def get_attachments(uid, folder, dest_dir):
"""Скачать вложения письма в dest_dir."""
"""
Скачать вложения письма в dest_dir.
Спек email-attachments: правильный флаг — `--downloads-dir` (не `--dir`).
Идемпотентность: если в dest_dir уже есть файлы — не качаем повторно.
"""
try:
existing = list(dest_dir.iterdir()) if dest_dir.exists() else []
if existing:
print(f" вложения уже скачаны ({len(existing)} ф.) — пропускаю")
return
except OSError:
pass
try:
run_cmd(
HIMALAYA_CMD + [
"attachment", "download", str(uid),
"--folder", folder,
"--dir", str(dest_dir),
"--downloads-dir", str(dest_dir),
],
timeout=60,
)
except RuntimeError:
pass # нет вложений — норм
except RuntimeError as e:
# Нет вложений / письмо не имеет вложений — норм для has_attachment=false.
# Но если письмо помечено has_attachment=true, а скачать не вышло —
# оставляем пустую папку и пишем warning (письмо не теряется).
print(f" [WARN] вложения не скачаны: {e}", file=sys.stderr)
def archive_folder(folder, limit=100):
@@ -454,10 +470,9 @@ def archive_folder(folder, limit=100):
date_str = env.get("date") or env.get("internal_date") or ""
year, month = parse_date(date_str)
# Путь: /mnt/yandex-disk/hermes/email/<folder>/YYYY/MM/UID/
# Путь: /opt/hermes/email/<folder>/YYYY/MM/UID/
msg_dir = ARCHIVE_ROOT / folder / f"{year:04d}" / f"{month:02d}" / str(uid)
email_path = msg_dir / "email.md"
attachments_dir = msg_dir / "attachments"
# Проверка — уже сохранено
if email_path.exists():
@@ -467,7 +482,13 @@ def archive_folder(folder, limit=100):
continue
msg_dir.mkdir(parents=True, exist_ok=True)
attachments_dir.mkdir(parents=True, exist_ok=True)
has_attachment = bool(env.get("has_attachment", False))
# Папку attachments/ создаём ТОЛЬКО если у письма есть вложения
# (спек email-attachments: без вложений пустую папку не создаём).
attachments_dir = None
if has_attachment:
attachments_dir = msg_dir / "attachments"
attachments_dir.mkdir(parents=True, exist_ok=True)
# Получаем заголовки и тело
extra_headers, body = get_email_content(uid, folder)
@@ -477,7 +498,8 @@ def archive_folder(folder, limit=100):
email_path.write_text(content, encoding="utf-8")
# Вложения
get_attachments(uid, folder, attachments_dir)
if attachments_dir is not None:
get_attachments(uid, folder, attachments_dir)
subj = (env.get("subject") or "")[:60]
print(f" ✓ UID {uid} ({subj})")