Files

648 lines
28 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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())