#!/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: """Читать сокет, пока не увидим строку ' 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) и дождаться ' 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 # дочитать до ' 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 :*). Вернуть список UID. ВАЖНО: UID — глобальный номер (сейчас ~14200+), а НЕ порядковый. Диапазонные запросы вида 'UID SEARCH UID 1:1000' пусты, если нет писем с такими UID (первый пустой батч обрывал цикл — письма помечались deleted). Поэтому используем целиковые запросы: * UID SEARCH ALL (работает, ~365 писем) * UID SEARCH UID :* (для новых, тоже работает) """ 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 # путь: /YYYY/MM//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())