mirror of
https://gitverse.ru/kpa39l/email-assistant.git
synced 2026-09-29 09:15:09 +00:00
Session 2026-09-15 code: imap_client (auth+metrics), handlers (TG date format), cron fix (venv python), openspec imap-realtime-sync proposal
This commit is contained in:
@@ -108,6 +108,17 @@ TG_TIMEOUT = 20
|
||||
# ---------------------------------------------------------------------------
|
||||
# Frontmatter
|
||||
# ---------------------------------------------------------------------------
|
||||
def _unquote_yaml(v):
|
||||
"""Снять обрамляющие двойные кавычки YAML-значения (артефакт простого парсера).
|
||||
|
||||
Если значение начинается/заканчивается на ", снять их и разэкранировать \\".
|
||||
"""
|
||||
v = v.strip()
|
||||
if len(v) >= 2 and v[0] == '"' and v[-1] == '"':
|
||||
v = v[1:-1].replace('\\"', '"').replace("\\\\", "\\")
|
||||
return v
|
||||
|
||||
|
||||
def parse_email_md(path):
|
||||
"""Прочитать email.md, вернуть (headers, body, content, fm_end)."""
|
||||
content = path.read_text(encoding="utf-8", errors="replace")
|
||||
@@ -120,7 +131,7 @@ def parse_email_md(path):
|
||||
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()
|
||||
headers[m.group(1)] = _unquote_yaml(m.group(2))
|
||||
else:
|
||||
headers, body, fm_end = {}, content.strip(), None
|
||||
return headers, body, content, fm_end
|
||||
@@ -379,15 +390,38 @@ def tg_call(method, **params):
|
||||
return data.get("result", {})
|
||||
|
||||
|
||||
def format_date_for_tg(date_str):
|
||||
"""'2026-09-02 10:49+03:00' → '02.09.2026 10:49' (ДД.ММ.ГГГГ ЧЧ:ММ).
|
||||
|
||||
Известные форматы frontmatter:
|
||||
'2026-09-02 10:49+03:00' (ISO, может быть пробел/буква перед таймзоной)
|
||||
'2026-09-02T10:49:03+03:00' (RFC3339)
|
||||
Если распарсить не удалось — вернуть как есть (не ломать уведомление).
|
||||
"""
|
||||
s = (date_str or "").strip()
|
||||
if not s:
|
||||
return ""
|
||||
t = s
|
||||
# нормализуем: отрезаем таймзону (+03:00 / +0300 / Z)
|
||||
m = re.match(r"^(\d{4})-(\d{2})-(\d{2})[T ](\d{2}):(\d{2})", s)
|
||||
if m:
|
||||
return f"{m.group(3)}.{m.group(2)}.{m.group(1)} {m.group(4)}:{m.group(5)}"
|
||||
return s
|
||||
|
||||
|
||||
def handle_urgent(path, headers, body):
|
||||
"""Отправить уведомление в Telegram. Возвращает (ok, detail)."""
|
||||
subject = (headers.get("subject") or "(без темы)").strip()
|
||||
sender = (headers.get("from") or "?").strip()
|
||||
# Дата/время получения — из frontmatter, формат ДД.ММ.ГГГГ ЧЧ:ММ
|
||||
date = format_date_for_tg(headers.get("date") or headers.get("internal_date"))
|
||||
date_line = f"📅 {date}\n" if date else ""
|
||||
preview = body.strip()
|
||||
if len(preview) > MAX_PREVIEW_CHARS:
|
||||
preview = preview[:MAX_PREVIEW_CHARS].rstrip() + "…"
|
||||
text = (
|
||||
f"⚠️ СРОЧНОЕ письмо\n\n"
|
||||
f"⚠️ СРОЧНОЕ письмо\n"
|
||||
f"{date_line}"
|
||||
f"От: {sender}\n"
|
||||
f"Тема: {subject}\n\n"
|
||||
f"{preview}\n\n"
|
||||
|
||||
@@ -0,0 +1,316 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
imap_client.py — общий IMAP-клиент для потоковой синхронизации почты.
|
||||
|
||||
Пока содержит вынесенную из mail_archive.py авторизацию как ОТДЕЛЬНУЮ функцию
|
||||
imap_connect(): соединение + STARTTLS + LOGIN с re-try (экспоненциальная пауза).
|
||||
|
||||
Зачем отдельный модуль:
|
||||
- mail_archive.py (опрос по cron) и imap_stream.py (IDLE-поток) будут
|
||||
использовать один и тот же клиент и одну авторизацию.
|
||||
- Авторизация на Microsoft Exchange (mail.corpoffice.tech) чувствительна к
|
||||
rate-limit: серия быстрых неудачных/частых подключений приводит к тому,
|
||||
что сервер начинает молчать (таймаут вместо OK/NO). Поэтому imap_connect()
|
||||
делает re-try с экспоненциальной паузой и ОДНОЙ точкой входа.
|
||||
|
||||
Пример:
|
||||
with imap_connect() as conn:
|
||||
conn.sendall(b"a2 LOGIN ...\r\n") # или просто команды поверх
|
||||
...
|
||||
|
||||
Возвращает (sock, creds). Сокет обёрнут в SSL (после STARTTLS).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import socket
|
||||
import ssl
|
||||
import time
|
||||
import sys
|
||||
import os
|
||||
import json
|
||||
import threading
|
||||
from pathlib import Path
|
||||
from datetime import datetime, timezone
|
||||
|
||||
# Краткая пауза между попытками: уважать rate-limit Exchange.
|
||||
RETRY_BASE = 5 # первая пауза, сек
|
||||
RETRY_MAX = 60 # потолок паузы
|
||||
RETRY_ATTEMPTS = 5
|
||||
|
||||
# ── Логирование и метрики ────────────────────────────────────────────────
|
||||
# Единый лог IMAP-клиента: /opt/hermes/email/logs/imap_client.log
|
||||
# Формат: JSON-строки (одна запись = одна строка), легко грепать/считать.
|
||||
LOG_DIR = Path("/opt/hermes/email/logs")
|
||||
LOG_FILE = LOG_DIR / "imap_client.log"
|
||||
_log_lock = threading.Lock()
|
||||
|
||||
|
||||
def _utcnow():
|
||||
return datetime.now(timezone.utc).isoformat(timespec="seconds")
|
||||
|
||||
|
||||
def imap_log(event, **fields):
|
||||
"""
|
||||
Записать событие в лог IMAP-клиента (JSON-строка).
|
||||
|
||||
Всегда пишутся: ts, event. Остальные поля — из kwargs.
|
||||
Пароль НИКОГДА не логируется (нет поля password нигде).
|
||||
Ошибка записи в лог не валит вызывающий код (try/except внутри).
|
||||
"""
|
||||
try:
|
||||
LOG_DIR.mkdir(parents=True, exist_ok=True)
|
||||
record = {"ts": _utcnow(), "event": event}
|
||||
record.update(fields)
|
||||
line = json.dumps(record, ensure_ascii=False)
|
||||
with _log_lock:
|
||||
with open(LOG_FILE, "a", encoding="utf-8") as f:
|
||||
f.write(line + "\n")
|
||||
except Exception:
|
||||
pass # лог не должен ломать IMAP
|
||||
|
||||
|
||||
def imap_metrics():
|
||||
"""
|
||||
Прочитать метрики из лога (доступность, успешность авторизации, сессии).
|
||||
|
||||
Возвращает dict:
|
||||
file — путь к логу
|
||||
total_auth — всего попыток авторизации (auth_ok + auth_failed)
|
||||
auth_ok — успешных
|
||||
auth_failed — неудачных
|
||||
auth_success_rate — доля успеха (0..1)
|
||||
conn_ok — соединений установлено
|
||||
conn_error — ошибок соединения/недоступности сервера
|
||||
last_session — последнее событие (ts, event, duration_ms, host)
|
||||
sessions_active — кол-во активных сессий (session_started - session_ended),
|
||||
по кратности событий в логе (приблизительно)
|
||||
"""
|
||||
try:
|
||||
rows = []
|
||||
if LOG_FILE.exists():
|
||||
with open(LOG_FILE, encoding="utf-8") as f:
|
||||
for line in f:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
rows.append(json.loads(line))
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
auth_ok = sum(1 for r in rows if r.get("event") == "auth_ok")
|
||||
auth_fail = sum(1 for r in rows if r.get("event") == "auth_failed")
|
||||
conn_ok = sum(1 for r in rows if r.get("event") == "conn_ok")
|
||||
conn_err = sum(1 for r in rows if r.get("event") == "conn_error")
|
||||
started = sum(1 for r in rows if r.get("event") == "session_started")
|
||||
ended = sum(1 for r in rows if r.get("event") == "session_ended")
|
||||
last = rows[-1] if rows else None
|
||||
return {
|
||||
"file": str(LOG_FILE),
|
||||
"total_auth": auth_ok + auth_fail,
|
||||
"auth_ok": auth_ok,
|
||||
"auth_failed": auth_fail,
|
||||
"auth_success_rate": round(auth_ok / (auth_ok + auth_fail), 3) if (auth_ok + auth_fail) else None,
|
||||
"conn_ok": conn_ok,
|
||||
"conn_error": conn_err,
|
||||
"last_event": last,
|
||||
"sessions_active": max(0, started - ended),
|
||||
"rows": len(rows),
|
||||
}
|
||||
except Exception as e:
|
||||
return {"error": str(e)}
|
||||
|
||||
|
||||
def _load_credentials():
|
||||
"""
|
||||
Достать IMAP-учётные данные из конфига himalaya (~/.config/himalaya/config.toml).
|
||||
Возвращает dict(host, port, login, password).
|
||||
(Вынесено из _himalaya_imap_credentials() в mail_archive.py, чтобы иметь
|
||||
единый источник кред и не дублировать парсинг конфига.)
|
||||
"""
|
||||
import tomllib
|
||||
from pathlib import Path
|
||||
|
||||
cfg_path = Path.home() / ".config" / "himalaya" / "config.toml"
|
||||
with open(cfg_path, "rb") as f:
|
||||
cfg = tomllib.load(f)
|
||||
accounts = cfg.get("accounts", {})
|
||||
name = next((n for n, a in accounts.items() if a.get("default")), None)
|
||||
if name is None and accounts:
|
||||
name = next(iter(accounts))
|
||||
if not name:
|
||||
raise RuntimeError("himalaya config: no account found")
|
||||
be = accounts[name].get("backend", {})
|
||||
auth = be.get("auth", {})
|
||||
password = auth.get("raw") or auth.get("password")
|
||||
if not password:
|
||||
cmd = auth.get("cmd", "")
|
||||
if cmd:
|
||||
# выполняем команду, выдающую пароль
|
||||
import subprocess
|
||||
password = subprocess.run(cmd.split(), capture_output=True,
|
||||
text=True, timeout=15).stdout.strip()
|
||||
if not password:
|
||||
raise RuntimeError(f"himalaya config: no password for account {name}")
|
||||
return {
|
||||
"host": be["host"],
|
||||
"port": be.get("port", 143),
|
||||
"login": be["login"],
|
||||
"password": password,
|
||||
}
|
||||
|
||||
|
||||
def imap_connect(host=None, port=None, login=None, password=None,
|
||||
retries=RETRY_ATTEMPTS, timeout=None):
|
||||
"""
|
||||
Установить IMAP-соединение с STARTTLS и авторизоваться (LOGIN).
|
||||
|
||||
Args:
|
||||
host/port/login/password: переопределение (по умолчанию — из конфига himalaya).
|
||||
retries: число попыток подключения+логина (для устойчивости к rate-limit).
|
||||
timeout: таймаут сокета, сек (по умолчанию 30, после логина 120 — Exchange
|
||||
медленно отвечает на большие FETCH).
|
||||
|
||||
Returns:
|
||||
(sock, creds): sock — SSL-сокет (после STARTTLS), creds — dict с host/login.
|
||||
|
||||
Raises:
|
||||
RuntimeError: если не удалось ни подключиться, ни авторизоваться.
|
||||
|
||||
Сокет НЕ закрывается при выходе — вызывающий закрывает (context manager
|
||||
в вызывающем коде). Авторизация выполняется простым LOGIN (как в работавшем
|
||||
fetch_attachments_imaplib): на этом Exchange он проходит, а AUTHENTICATE PLAIN
|
||||
отклоняется (rate-limit/fingerprint — см. README/design).
|
||||
"""
|
||||
if host is None:
|
||||
creds = _load_credentials()
|
||||
host, port, login, password = creds["host"], creds["port"], creds["login"], creds["password"]
|
||||
else:
|
||||
creds = {"host": host, "port": port, "login": login, "password": password}
|
||||
|
||||
if timeout is None:
|
||||
timeout = 30
|
||||
|
||||
last_err = None
|
||||
for attempt in range(1, retries + 1):
|
||||
sock = None
|
||||
t_start = time.monotonic()
|
||||
try:
|
||||
sock = socket.create_connection((host, port), timeout=timeout)
|
||||
sock.settimeout(timeout)
|
||||
greet = sock.recv(1024)
|
||||
if not greet.startswith(b"* OK"):
|
||||
raise RuntimeError(f"bad greeting: {greet[:80]!r}")
|
||||
sock.sendall(b"a1 STARTTLS\r\n")
|
||||
resp = sock.recv(1024)
|
||||
if b"OK" not in resp:
|
||||
raise RuntimeError(f"STARTTLS failed: {resp[:80]!r}")
|
||||
ctx = ssl.create_default_context()
|
||||
sock = ctx.wrap_socket(sock, server_hostname=host)
|
||||
sock.settimeout(120) # после авторизации — на большие FETCH
|
||||
|
||||
# LOGIN — простой, как в fetch_attachments_imaplib (работает на Exchange)
|
||||
login_cmd = f'a2 LOGIN {login} {password}\r\n'.encode()
|
||||
sock.sendall(login_cmd)
|
||||
buf = b""
|
||||
while True:
|
||||
d = sock.recv(65536)
|
||||
if not d:
|
||||
break
|
||||
buf += d
|
||||
# дождаться строки "a2 OK/NO/BAD"
|
||||
lines = buf.split(b"\r\n")
|
||||
if any(l.startswith(b"a2 ") for l in lines):
|
||||
break
|
||||
if b"a2 OK" not in buf:
|
||||
# последняя строка (ответ) для диагностики, пароль не показываем
|
||||
tail = buf[-200:].decode("utf-8", "replace")
|
||||
imap_log("auth_failed", host=host, port=port, login=login,
|
||||
attempt=attempt, reason="LOGIN rejected",
|
||||
detail=tail, duration_ms=int((time.monotonic() - t_start) * 1000))
|
||||
raise RuntimeError(f"LOGIN failed: ...{tail}")
|
||||
|
||||
# Успех: метрика доступности (conn_ok) + успешная авторизация
|
||||
imap_log("conn_ok", host=host, port=port, attempt=attempt,
|
||||
duration_ms=int((time.monotonic() - t_start) * 1000))
|
||||
imap_log("auth_ok", host=host, login=login, attempt=attempt,
|
||||
duration_ms=int((time.monotonic() - t_start) * 1000))
|
||||
imap_log("session_started", host=host, login=login)
|
||||
|
||||
# вернуть сокет; пароль НЕ логируем
|
||||
return sock, creds
|
||||
|
||||
except (socket.timeout, ConnectionError, ssl.SSLError, RuntimeError) as e:
|
||||
last_err = e
|
||||
# Классифицируем для метрики доступности: ошибка соединения/таймаут =
|
||||
# проблема доступности сервера; обычный RuntimeError = проблема
|
||||
# авторизации/протокола (уже залогирован как auth_failed выше).
|
||||
if isinstance(e, (socket.timeout, ConnectionError, ssl.SSLError)):
|
||||
imap_log("conn_error", host=host, port=port, attempt=attempt,
|
||||
error=type(e).__name__,
|
||||
duration_ms=int((time.monotonic() - t_start) * 1000),
|
||||
detail=str(e)[:200])
|
||||
if sock is not None:
|
||||
try:
|
||||
sock.close()
|
||||
except Exception:
|
||||
pass
|
||||
if attempt < retries:
|
||||
pause = min(RETRY_BASE * (2 ** (attempt - 1)), RETRY_MAX)
|
||||
time.sleep(pause)
|
||||
|
||||
imap_log("auth_failed", host=host, port=port, login=login,
|
||||
attempt=retries, reason=f"all attempts failed: {type(last_err).__name__}: {last_err}"[:200])
|
||||
raise RuntimeError(f"IMAP connect/login failed after {retries} attempts: {last_err}")
|
||||
|
||||
|
||||
class imap_session:
|
||||
"""
|
||||
Контекстный менеджер IMAP-сессии: авторизация + гарантированная запись
|
||||
session_ended при закрытии (для метрики "активная сессия").
|
||||
|
||||
with imap_session() as (sock, creds):
|
||||
sock.sendall(...)
|
||||
|
||||
Эквивалентно imap_connect(), но на выходе пишет в лог session_ended
|
||||
(даже при исключении). Пароль не логируется.
|
||||
"""
|
||||
|
||||
def __init__(self, **kwargs):
|
||||
self._kwargs = kwargs
|
||||
self._sock = None
|
||||
self._creds = None
|
||||
|
||||
def __enter__(self):
|
||||
self._sock, self._creds = imap_connect(**self._kwargs)
|
||||
return self._sock, self._creds
|
||||
|
||||
def __exit__(self, exc_type, exc, tb):
|
||||
if self._sock is not None:
|
||||
try:
|
||||
imap_log("session_ended", host=self._creds.get("host") if self._creds else None,
|
||||
login=self._creds.get("login") if self._creds else None)
|
||||
finally:
|
||||
try:
|
||||
self._sock.close()
|
||||
except Exception:
|
||||
pass
|
||||
self._sock = None
|
||||
return False # не глотаем исключения
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
# Быстрый самотест: установить соединение и закрыть.
|
||||
# python3 imap_client.py — проверить авторизацию
|
||||
# python3 imap_client.py --metrics — показать метрики из лога
|
||||
if len(sys.argv) > 1 and sys.argv[1] == "--metrics":
|
||||
m = imap_metrics()
|
||||
print(json.dumps(m, ensure_ascii=False, indent=2))
|
||||
sys.exit(0)
|
||||
try:
|
||||
with imap_session() as (s, creds):
|
||||
print(f"OK: connected+LOGIN as {creds['login']} @ {creds['host']}")
|
||||
s.sendall(b"a3 LOGOUT\r\n")
|
||||
except Exception as e:
|
||||
print(f"FAIL: {e}")
|
||||
sys.exit(1)
|
||||
@@ -5,8 +5,14 @@ set -euo pipefail
|
||||
|
||||
cd /opt/hermes/email-assistant
|
||||
|
||||
# Классификатор: до 10 новых писем за запуск (Qwen ~10-20с/письмо = ~3 мин)
|
||||
python3 scripts/email_classifier.py --limit 10 || echo "[classifier] ошибка (продолжаем)" >&2
|
||||
# httpx (для Telegram в email_handlers.py) есть только в .venv /opt/vesti.
|
||||
# Системный python3 его не имеет — без этого уведомления не отправлялись.
|
||||
PY=/opt/vesti/.venv/bin/python
|
||||
[ -x "$PY" ] || PY=python3
|
||||
|
||||
# Обработчики: все письма с тегами, без handled_* (идемпотентно)
|
||||
python3 scripts/email_handlers.py || echo "[handlers] ошибка" >&2
|
||||
# Классификатор: до 10 новых писем за запуск (Qwen ~10-20с/письмо = ~3 мин)
|
||||
"$PY" scripts/email_classifier.py --limit 10 || echo "[classifier] ошибка (продолжаем)" >&2
|
||||
|
||||
# Обработчики: по 2 письма за проход (идемпотентно; избегаем залпа
|
||||
# уведомлений, если накопился backlog urgent)
|
||||
"$PY" scripts/email_handlers.py --limit 2 || echo "[handlers] ошибка" >&2
|
||||
+9
-22
@@ -480,26 +480,16 @@ def fetch_attachments_imaplib(uid, folder, dest_dir):
|
||||
import email as emailmod
|
||||
import email.header as email_header
|
||||
|
||||
try:
|
||||
creds = _himalaya_imap_credentials()
|
||||
except Exception as e:
|
||||
print(f" [WARN] нет IMAP-учётных из конфига himalaya: {e}", file=sys.stderr)
|
||||
return False
|
||||
|
||||
sock = None
|
||||
try:
|
||||
sock = socket.create_connection((creds["host"], creds["port"]), timeout=30)
|
||||
sock.settimeout(30)
|
||||
greet = sock.recv(1024)
|
||||
if not greet.startswith(b"* OK"):
|
||||
raise RuntimeError(f"bad greeting: {greet[:80]!r}")
|
||||
sock.sendall(b"a1 STARTTLS\r\n")
|
||||
resp = sock.recv(1024)
|
||||
if b"OK" not in resp:
|
||||
raise RuntimeError(f"STARTTLS failed: {resp[:80]!r}")
|
||||
ctx = sslmod.create_default_context()
|
||||
sock = ctx.wrap_socket(sock, server_hostname=creds["host"])
|
||||
sock.settimeout(30)
|
||||
# Авторизация — через общий imap_client.imap_connect() (отдельная функция,
|
||||
# с re-try против rate-limit Exchange). Креды читаются из конфига himalaya
|
||||
# внутри imap_connect(). Возвращает (sock, creds).
|
||||
import sys as _sys
|
||||
from pathlib import Path as _Path
|
||||
_sys.path.insert(0, str(_Path(__file__).resolve().parent))
|
||||
from imap_client import imap_connect
|
||||
sock, creds = imap_connect()
|
||||
|
||||
def cmd(tag, line, expect_literal=None):
|
||||
sock.sendall(f"{tag} {line}\r\n".encode())
|
||||
@@ -525,9 +515,6 @@ def fetch_attachments_imaplib(uid, folder, dest_dir):
|
||||
break
|
||||
return buf
|
||||
|
||||
r = cmd("a2", f'LOGIN {creds["login"]} {creds["password"]}')
|
||||
if b"OK LOGIN" not in r:
|
||||
raise RuntimeError(f"LOGIN failed: {r[-120:]!r}")
|
||||
# Папка — в IMAP modified UTF-7 (кириллица иначе не находится на Exchange)
|
||||
mbox = _imap_utf7_encode(folder)
|
||||
r = cmd("a3", f'SELECT "{mbox}"')
|
||||
@@ -589,7 +576,7 @@ def fetch_attachments_imaplib(uid, folder, dest_dir):
|
||||
pass
|
||||
|
||||
|
||||
def get_attachments(uid, folder, dest_dir, timeout=25):
|
||||
def get_attachments(uid, folder, dest_dir, timeout=120):
|
||||
"""
|
||||
Скачать вложения письма в dest_dir.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user