imap_stream (Фаза 1): IDLE-цикл с реконнектом + SQLite mailbox.db; фикс токена TG (vesti/.env всегда); фильтр актуальности urgent (>=3 дн пропуск); HERMES_REAL_HOME в imap_client

This commit is contained in:
2026-09-15 18:34:37 +00:00
parent defcf53d75
commit 5a81c86599
3 changed files with 692 additions and 3 deletions
+647
View File
@@ -0,0 +1,647 @@
#!/usr/bin/env python3
"""
imap_stream.py — постоянный IMAP IDLE-поток (Фаза 1 imap-realtime-sync).
Единственный процесс, который держит постоянное соединение с IMAP
(mail.corpoffice.tech, Microsoft Exchange) и видит события почтового ящика
мгновенно, как почтовый клиент:
* N EXISTS — новые письма
* N EXPUNGE — удаления (порядковый номер; UID выясняем reconcile-ом)
* N FETCH FLAGS — смена флагов (прочитано/ответ/флаг)
Пишет события в SQLite ChangeLog (/opt/hermes/email/state/mailbox.db):
mailbox_state — онлайновая копия состояния ящика (uid, folder, flags,
message_id, deleted)
mailbox_events — append-only журнал изменений (added/moved/deleted/
flag_changed/reconcile)
Новые письма архивируются переиспользованием mail_archive.py
(himalaya --preview + BODY.PEEK[] — флаг \\Seen НЕ ставится).
Цикл:
connect (imap_client.imap_connect, re-try) → SELECT INBOX → IDLE
→ на событие или таймаут (~28 мин, сервер рвёт ~30): DONE
→ быстрый reconcile (новые UID) + периодический полный (флаги/удаления)
→ снова IDLE. При обрыве: пауза >= CONNECT_PAUSE (rate-limit Exchange)
→ reconnect + полный reconcile.
Запуск:
python3 imap_stream.py # демон (бесконечный цикл)
python3 imap_stream.py --check # connect+SELECT+LOGOUT, выйти
python3 imap_stream.py --test-idle# connect+SELECT+IDLE 30с+DONE+reconcile, выйти
python3 imap_stream.py --status # краткий статус ChangeLog/state из БД
\\Seen никогда не ставится: stream не делает BODY[] — только SEARCH/FETCH FLAGS;
архивация — через himalaya --preview / BODY.PEEK[] (см. mail_archive.py).
"""
from __future__ import annotations
import argparse
import json
import os
import re
import signal
import socket
import sqlite3
import sys
import time
from datetime import datetime, timezone
from pathlib import Path
# Общий IMAP-клиент (auth + метрики) — рядом с нами в scripts/
sys.path.insert(0, str(Path(__file__).resolve().parent))
from imap_client import imap_connect, imap_log # noqa: E402
from mail_archive import ARCHIVE_ROOT, archive_folder, get_inbox_subfolders # noqa: E402
# ── Конфигурация ──────────────────────────────────────────────────────────
STATE_DIR = Path("/opt/hermes/email/state")
DB_PATH = STATE_DIR / "mailbox.db"
INBOX = "INBOX"
# Rate-limit Exchange: пауза между ПОДКЛЮЧЕНИЯМИ (не командами). После серии
# быстрых коннектов сервер начинает молчать (timeout). Ниже 45с не опускаться
# (решение сессии 2026-09-15, live-тесты).
CONNECT_PAUSE = 45 # сек, пауза после обрыва перед reconnect
# Exchange рвёт IDLE-соединение ~60с (проверено live 2026-09-15), поэтому
# перевыпускаем IDLE с запасом — 25с (до серверного лимита).
IDLE_TIMEOUT = 25 # сек, перевыпуск IDLE (сервер рвёт ~60с)
RECONCILE_FULL_EVERY = 5 # каждый N-й цикл IDLE — полный reconcile (флаги/удаления)
FETCH_FLAGS_CHUNK = 500 # UID FETCH (FLAGS) батчами по 500
# Письма новее last_uid архивируем через mail_archive.archive_folder
ARCHIVE_LIMIT = 200 # максимум писем за один вызов архиватора
# ── SQLite: mailbox_state + mailbox_events ────────────────────────────────
SCHEMA = """
CREATE TABLE IF NOT EXISTS mailbox_state (
uid INTEGER NOT NULL,
folder TEXT NOT NULL,
message_id TEXT,
in_reply_to TEXT,
refs TEXT,
flags TEXT,
has_attachment INTEGER DEFAULT 0,
archive_path TEXT,
last_seen TEXT,
deleted INTEGER DEFAULT 0,
PRIMARY KEY (folder, uid)
);
CREATE TABLE IF NOT EXISTS mailbox_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ts TEXT NOT NULL,
event TEXT NOT NULL, -- added|moved|deleted|flag_changed|replied|reconcile
folder TEXT,
uid INTEGER,
message_id TEXT,
details TEXT
);
CREATE INDEX IF NOT EXISTS idx_events_message ON mailbox_events(message_id);
CREATE INDEX IF NOT EXISTS idx_events_ts ON mailbox_events(ts);
"""
def _utcnow() -> str:
return datetime.now(timezone.utc).isoformat(timespec="seconds")
def init_db(db_path: Path = DB_PATH):
STATE_DIR.mkdir(parents=True, exist_ok=True)
conn = sqlite3.connect(str(db_path))
conn.executescript(SCHEMA)
conn.commit()
return conn
def insert_event(conn, event, folder=None, uid=None, message_id=None, details=None):
cur = conn.execute(
"INSERT INTO mailbox_events (ts, event, folder, uid, message_id, details)"
" VALUES (?, ?, ?, ?, ?, ?)",
(_utcnow(), event, folder, uid, message_id, details),
)
conn.commit()
return cur.lastrowid
def upsert_state(conn, folder, uid, *, message_id=None, in_reply_to=None,
references=None, flags=None, has_attachment=None,
archive_path=None, deleted=None):
"""Обновить/вставить строку mailbox_state (PK folder+uid)."""
cur = conn.execute(
"SELECT flags, deleted FROM mailbox_state WHERE folder=? AND uid=?",
(folder, uid),
)
row = cur.fetchone()
if row is None:
conn.execute(
"INSERT INTO mailbox_state (uid, folder, message_id, in_reply_to,"
" refs, flags, has_attachment, archive_path, last_seen, deleted)"
" VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 0)",
(uid, folder, message_id, in_reply_to, references, flags,
has_attachment or 0, archive_path, _utcnow()),
)
conn.commit()
return "inserted"
new_flags = flags if flags is not None else row[0]
new_deleted = deleted if deleted is not None else row[1]
conn.execute(
"UPDATE mailbox_state SET flags=?, deleted=?, last_seen=? WHERE folder=? AND uid=?",
(new_flags, new_deleted, _utcnow(), folder, uid),
)
conn.commit()
return "updated"
def state_uids(conn, folder=INBOX):
cur = conn.execute(
"SELECT uid, flags, message_id FROM mailbox_state WHERE folder=? AND deleted=0",
(folder,),
)
return {row[0]: {"flags": row[1], "message_id": row[2]} for row in cur.fetchall()}
def load_last_uid(folder=INBOX):
"""Последний заархивированный UID — из state-файла mail_archive (общий с cron)."""
sf = STATE_DIR / f"mail-archive-last-{folder}.json"
if sf.exists():
try:
return int(json.loads(sf.read_text()).get("last_uid", 0))
except (ValueError, json.JSONDecodeError):
return 0
return 0
def load_mailbox_last_uid(conn, folder=INBOX):
cur = conn.execute("SELECT MAX(uid) FROM mailbox_state WHERE folder=?", (folder,))
return cur.fetchone()[0] or 0
# ── IMAP-обёртка поверх сокета (сырой, как imap_client) ──────────────────
class ImapStream:
def __init__(self, conn, creds):
self.sock, self.creds = conn
self.sock.settimeout(120)
self._tag = 0
self._buf = b""
self._selected = None
# -- низкоуровневые команды -------------------------------------------
def _next_tag(self):
self._tag += 1
return f"a{self._tag}"
def _read_until_tag(self, tag: str, timeout: float | None = None) -> bytes:
"""Читать сокет, пока не увидим строку '<tag> OK/NO/BAD' (или таймаут)."""
if timeout is not None:
self.sock.settimeout(timeout)
else:
self.sock.settimeout(120)
end_marker = tag.encode() + b" "
while True:
if end_marker in self._buf:
break
try:
d = self.sock.recv(65536)
except socket.timeout:
raise
if not d:
raise ConnectionError("connection closed by server")
self._buf += d
# отделить ответ от буфера
idx = self._buf.find(end_marker)
resp = self._buf[: idx + len(end_marker)]
# дочитать строку до \r\n (хвост ответа тега)
while b"\r\n" not in self._buf[idx:]:
d = self.sock.recv(65536)
if not d:
break
self._buf += d
nl = self._buf.find(b"\r\n", idx)
if nl != -1:
resp += self._buf[idx + len(end_marker): nl + 2]
self._buf = self._buf[nl + 2:]
else:
self._buf = b""
return resp
def _cmd(self, line: str, timeout: float | None = None) -> bytes:
"""Отправить команду, вернуть весь ответ (untagged + tagged)."""
tag = self._next_tag()
self.sock.sendall(f"{tag} {line}\r\n".encode())
resp = self._read_until_tag(tag, timeout=timeout)
return resp
def _cmd_ok(self, line: str, timeout: float | None = None) -> bytes:
"""Команда, требующая 'tag OK'. RuntimeError при NO/BAD."""
resp = self._cmd(line, timeout=timeout)
# последняя непустая строка — ответ с тегом; после \r\n бывает пустой хвост
lines = [l for l in resp.split(b"\r\n") if l.strip()]
last = lines[-1] if lines else b""
if not last.startswith(b"a") or b" OK " not in last:
tail = resp[-160:].decode("utf-8", "replace")
raise RuntimeError(f"IMAP command failed: {line!r} -> ...{tail}")
return resp
# -- высокоуровневые операции -----------------------------------------
def select(self, folder=INBOX):
resp = self._cmd_ok(f"SELECT \"{folder}\"")
self._selected = folder
# вернём число существующих писем (* N EXISTS)
m = re.search(rb"\* (\d+) EXISTS", resp)
return int(m.group(1)) if m else 0
def logout(self):
try:
self.sock.sendall(b"LOGOUT\r\n")
except Exception:
pass
try:
self.sock.close()
except Exception:
pass
# -- IDLE --------------------------------------------------------------
def idle_start(self):
"""Начать IDLE. Возвращает True, если сервер ответил '+ ...'.
Exchange отвечает '+ IDLE accepted, awaiting DONE command.'
(не '+ idling') — матчим любой '+ ' в начале строки.
"""
tag = self._next_tag()
self.sock.sendall(f"{tag} IDLE\r\n".encode())
# ждём '+ ' (положительный ответ сервера)
deadline = time.monotonic() + 30
while b"+ IDLE" not in self._buf and b"+ idling" not in self._buf:
if time.monotonic() > deadline:
raise RuntimeError("IDLE: no + from server (timeout 30s)")
self.sock.settimeout(35)
d = self.sock.recv(65536)
if not d:
raise ConnectionError("connection closed during IDLE start")
self._buf += d
# вычистить строку '+ idling'
nl = self._buf.find(b"\r\n")
if nl != -1:
self._buf = self._buf[nl + 2:]
self._idle_tag = tag
return True
def idle_wait(self, timeout: float) -> list[bytes]:
"""Ждать untagged-события до timeout сек. Вернуть строки событий.
На таймаут — вернуть [] (это нормальный момент перевыпуска IDLE).
На событие — вернуть накопленные untagged-строки, IDLE продолжает висеть.
На обрыв/ошибку — raise.
"""
events: list[bytes] = []
self.sock.settimeout(timeout)
try:
while True:
d = self.sock.recv(65536)
if not d:
raise ConnectionError("connection closed during IDLE")
self._buf += d
# события — строки '* ...' в буфере (до \r\n)
while b"\r\n" in self._buf:
line, self._buf = self._buf.split(b"\r\n", 1)
line = line.strip()
if not line:
continue
if line.startswith(b"*"):
events.append(line)
except socket.timeout:
return events
def idle_done(self):
"""Завершить IDLE (DONE) и дождаться '<tag> OK IDLE completed'."""
if not getattr(self, "_idle_tag", None):
return
tag = self._idle_tag
if tag is None:
return
try:
self.sock.sendall(b"DONE\r\n")
except Exception:
return
# дочитать до '<tag> OK'
end_marker = tag.encode() + b" "
deadline = time.monotonic() + 30
while end_marker not in self._buf and time.monotonic() < deadline:
try:
d = self.sock.recv(65536)
except socket.timeout:
break
if not d:
break
self._buf += d
self._idle_tag = None
# вычистить остатки ответа DONE
if end_marker in self._buf:
idx = self._buf.find(end_marker)
nl = self._buf.find(b"\r\n", idx)
if nl != -1:
self._buf = self._buf[nl + 2:]
# -- поиск и флаги -----------------------------------------------------
def search_uids(self, folder=INBOX, uid_min=None) -> list[int]:
"""UID SEARCH ALL (или UID <min>:*). Вернуть список UID.
ВАЖНО: UID — глобальный номер (сейчас ~14200+), а НЕ порядковый.
Диапазонные запросы вида 'UID SEARCH UID 1:1000' пусты, если нет
писем с такими UID (первый пустой батч обрывал цикл — письма
помечались deleted). Поэтому используем целиковые запросы:
* UID SEARCH ALL (работает, ~365 писем)
* UID SEARCH UID <min>:* (для новых, тоже работает)
"""
if uid_min is not None:
resp = self._cmd(f"UID SEARCH UID {uid_min}:*")
else:
resp = self._cmd("UID SEARCH ALL")
m = re.search(rb"\* SEARCH(.*?)\r\n", resp)
if not m:
return []
return [int(x) for x in m.group(1).split()]
def fetch_flags(self, uids: list[int], folder=INBOX) -> dict[int, str]:
"""UID FETCH (FLAGS) батчами. Вернуть {uid: 'FLAGS-строка'}."""
out: dict[int, str] = {}
for i in range(0, len(uids), FETCH_FLAGS_CHUNK):
chunk = uids[i:i + FETCH_FLAGS_CHUNK]
uid_list = ",".join(str(u) for u in chunk)
resp = self._cmd(f"UID FETCH {uid_list} (FLAGS)")
# строки вида: * 123 FETCH (FLAGS (\\Seen) UID 456)
for line in resp.split(b"\r\n"):
if not line.startswith(b"*"):
continue
m = re.search(rb"FLAGS \(([^)]*)\)", line)
mu = re.search(rb"UID (\d+)", line)
if m and mu:
out[int(mu.group(1))] = m.group(1).decode("utf-8", "replace")
return out
# -- reconcile ---------------------------------------------------------
def reconcile_new(self, conn, folder=INBOX):
"""Найти письма UID > last_uid, записать added + заархивировать.
Возвращает список новых UID.
"""
last_uid = max(load_last_uid(folder), load_mailbox_last_uid(conn, folder))
new_uids = self.search_uids(folder, uid_min=last_uid + 1)
new_uids = [u for u in new_uids if u > last_uid]
if not new_uids:
return []
for uid in new_uids:
insert_event(conn, "added", folder=folder, uid=uid,
details=f"detected by reconcile (uid>{last_uid})")
upsert_state(conn, folder, uid, flags="")
conn.commit()
imap_log("stream_reconcile", folder=folder, new=len(new_uids), last_uid=last_uid)
# Архивация — переиспользуем mail_archive (himalaya --preview, \\Seen не ставит).
# Он сам обновит state-файл last_uid (общий с cron, идемпотентно).
try:
archived = archive_folder(folder, limit=ARCHIVE_LIMIT)
imap_log("stream_archived", folder=folder, count=archived)
except Exception as e:
imap_log("stream_archived_error", folder=folder, error=str(e)[:200])
# Обновить message_id в state из заархивированных email.md
self._fill_message_ids(conn, folder, new_uids)
return new_uids
def _fill_message_ids(self, conn, folder, uids):
"""Для каждого UID прочитать Message-ID из email.md (frontmatter)."""
for uid in uids:
path = self._find_email_md(folder, uid)
if not path:
continue
try:
text = path.read_text(encoding="utf-8", errors="ignore")
mid = re.search(r"(?m)^Message-ID:\s*(.+)$", text)
in_reply = re.search(r"(?m)^In-Reply-To:\s*(.+)$", text)
refs = re.search(r"(?m)^References:\s*(.+)$", text)
conn.execute(
"UPDATE mailbox_state SET message_id=?, in_reply_to=?,"
" refs=?, archive_path=?, has_attachment=? WHERE folder=? AND uid=?",
(mid.group(1).strip() if mid else None,
in_reply.group(1).strip() if in_reply else None,
refs.group(1).strip() if refs else None,
str(path.parent),
1 if (path.parent / "attachments").exists() else 0,
folder, uid),
)
except Exception:
pass
conn.commit()
def _find_email_md(self, folder, uid):
base = ARCHIVE_ROOT / folder
if not base.exists():
return None
# путь: <folder>/YYYY/MM/<uid>/email.md
hits = list(base.glob(f"*/*/{uid}/email.md"))
return hits[0] if hits else None
def reconcile_full(self, conn, folder=INBOX):
"""Полный reconcile: флаги известных писем + deleted для пропавших.
Делается периодически (IDLE не гарантирует доставку всех событий;
после обрыва — обязательно).
"""
known = state_uids(conn, folder)
if not known:
return
all_uids = set(self.search_uids(folder))
# пропавшие (удалены на сервере) → deleted (soft)
missing = [u for u in known if u not in all_uids]
for uid in missing:
insert_event(conn, "deleted", folder=folder, uid=uid,
message_id=known[uid]["message_id"],
details="uid missing on server (expunged)")
conn.execute("UPDATE mailbox_state SET deleted=1 WHERE folder=? AND uid=?",
(folder, uid))
conn.commit()
# флаги изменились → flag_changed
changed_flags = self.fetch_flags(sorted(known), folder)
for uid, flags in changed_flags.items():
prev = known.get(uid)
if prev is None:
continue
prev_flags = prev["flags"] or ""
if prev_flags != flags:
insert_event(conn, "flag_changed", folder=folder, uid=uid,
message_id=prev["message_id"],
details=f"{prev_flags or '()'} -> {flags or '()'}")
conn.execute("UPDATE mailbox_state SET flags=? WHERE folder=? AND uid=?",
(flags, folder, uid))
conn.commit()
if missing or changed_flags:
imap_log("stream_reconcile_full", folder=folder,
deleted=len(missing), flags_changed=len(changed_flags))
return len(missing), len(changed_flags)
# -- обработка untagged IDLE-событий -----------------------------------
def handle_idle_events(self, events: list[bytes]) -> str:
"""Классифицировать untagged-строки IDLE. Вернуть тип события."""
kinds = set()
for line in events:
if b"EXISTS" in line:
kinds.add("exists")
elif b"EXPUNGE" in line:
kinds.add("expunge")
elif b"FETCH" in line and b"FLAGS" in line:
kinds.add("flags")
imap_log("stream_idle_event", events=[e.decode("utf-8", "replace") for e in events])
# приоритет: expunge > exists > flags (expunge требует полного reconcile)
if "expunge" in kinds:
return "expunge"
if "exists" in kinds:
return "exists"
if "flags" in kinds:
return "flags"
return "none"
# -- главный цикл ------------------------------------------------------
def run(self, db_path=DB_PATH):
conn = init_db(db_path)
cycle = 0
imap_log("stream_started", host=self.creds["host"], login=self.creds["login"])
while True:
cycle += 1
try:
exists = self.select(INBOX)
imap_log("stream_select", folder=INBOX, exists=exists)
# Полный reconcile на старте и после обрыва (IDLE не гарантирует события)
self.reconcile_full(conn)
self.reconcile_new(conn)
while True:
self.idle_start()
events = self.idle_wait(IDLE_TIMEOUT)
kind = self.handle_idle_events(events) if events else "timeout"
self.idle_done()
if kind == "timeout":
# перевыпуск IDLE + лёгкий reconcile новых
self.reconcile_new(conn)
if cycle % RECONCILE_FULL_EVERY == 0:
self.reconcile_full(conn)
imap_log("stream_idle_reissue", cycle=cycle)
continue
# событие: быстрый reconcile новых + полный (для expunge)
self.reconcile_new(conn)
if kind in ("expunge", "flags") or cycle % RECONCILE_FULL_EVERY == 0:
self.reconcile_full(conn)
imap_log("stream_event_processed", kind=kind, cycle=cycle)
except (ConnectionError, socket.timeout, OSError, RuntimeError) as e:
imap_log("stream_error", error=type(e).__name__, detail=str(e)[:200])
try:
self.logout()
except Exception:
pass
# Rate-limit Exchange: пауза между коннектами >= 45с
imap_log("stream_reconnect_pause", seconds=CONNECT_PAUSE)
time.sleep(CONNECT_PAUSE)
try:
self.sock, self.creds = imap_connect()
except RuntimeError as ce:
imap_log("stream_reconnect_failed", error=str(ce)[:200])
time.sleep(CONNECT_PAUSE * 2)
# переподключение на следующей итерации цикла
continue
def status(db_path=DB_PATH):
if not db_path.exists():
print("mailbox.db ещё не создан")
return
conn = sqlite3.connect(str(db_path))
total = conn.execute("SELECT COUNT(*) FROM mailbox_state").fetchone()[0]
deleted = conn.execute("SELECT COUNT(*) FROM mailbox_state WHERE deleted=1").fetchone()[0]
ev = conn.execute("SELECT event, COUNT(*) FROM mailbox_events GROUP BY event ORDER BY 2 DESC").fetchall()
print(f"mailbox_state: {total} строк (deleted={deleted})")
print("mailbox_events:")
for row in ev:
print(f" {row[0]}: {row[1]}")
last = conn.execute("SELECT id, ts, event, folder, uid FROM mailbox_events ORDER BY id DESC LIMIT 5").fetchall()
print("последние события:")
for row in last:
print(f" #{row[0]} {row[1]} {row[2]} {row[3]} uid={row[4]}")
def main():
ap = argparse.ArgumentParser(description="IMAP IDLE-поток (реалтайм синк почты)")
ap.add_argument("--check", action="store_true",
help="connect + SELECT + LOGOUT, выйти (проверка авторизации)")
ap.add_argument("--test-idle", action="store_true",
help="connect + SELECT + IDLE 30с + DONE + reconcile, выйти")
ap.add_argument("--status", action="store_true", help="статус mailbox.db")
args = ap.parse_args()
if args.status:
status()
return 0
if args.check:
with imap_session_ctx() as (sock, creds):
print(f"OK: connected+LOGIN as {creds['login']} @ {creds['host']}")
st = ImapStream((sock, creds), creds)
n = st.select(INBOX)
print(f"OK: SELECT {INBOX} — {n} писем")
st.logout()
return 0
if args.test_idle:
# тест: один цикл IDLE 30с + reconcile, потом выход
sock, creds = imap_connect()
st = ImapStream((sock, creds), creds)
conn = init_db()
try:
n = st.select(INBOX)
print(f"SELECT {INBOX}: {n} писем")
print("reconcile_new...")
new = st.reconcile_new(conn)
print(f" новых: {len(new)} {new[:10]}")
print("reconcile_full...")
d, fc = st.reconcile_full(conn)
print(f" deleted={d or 0}, flags_changed={fc or 0}")
if d or fc:
print(" изменения: см. mailbox_events")
print("IDLE 30с (жду событий, потом DONE)...")
st.idle_start()
ev = st.idle_wait(30)
print(f" событий за 30с: {len(ev)}")
for e in ev:
print(f" {e.decode('utf-8', 'replace')}")
st.idle_done()
print("IDLE завершён OK")
finally:
st.logout()
return 0
# демон
sock, creds = imap_connect()
st = ImapStream((sock, creds), creds)
def _sigterm(signum, frame):
imap_log("stream_stopped", signal=signum)
sys.exit(0)
signal.signal(signal.SIGTERM, _sigterm)
signal.signal(signal.SIGINT, _sigterm)
try:
st.run()
except KeyboardInterrupt:
imap_log("stream_stopped", signal="SIGINT")
return 0
def imap_session_ctx():
"""Мини-контекст для --check (session_ended в лог)."""
from imap_client import imap_session
return imap_session()
if __name__ == "__main__":
sys.exit(main())