diff --git a/openspec/changes/imap-realtime-sync/design.md b/openspec/changes/imap-realtime-sync/design.md new file mode 100644 index 0000000..fe893f7 --- /dev/null +++ b/openspec/changes/imap-realtime-sync/design.md @@ -0,0 +1,212 @@ +# Design: Потоковая синхронизация почты (IMAP IDLE) и онлайновая копия ящика + +## Context + +- Сервер: `mail.corpoffice.tech:143` (STARTTLS), Microsoft Exchange. + Проверено 2026-09-15 (ручной IMAP-тест): + - CAPABILITY: `IMAP4 IMAP4rev1 AUTH=PLAIN AUTH=NTLM AUTH=GSSAPI STARTTLS + SASL-IR UIDPLUS MOVE ID UNSELECT CHILDREN IDLE NAMESPACE LITERAL+` + - **IDLE, MOVE, UIDPLUS, CHILDREN, LITERAL+ — всё есть**, значит потоковая + синхронизация и отслеживание перемещений возможны. + - Логин с чистого скрипта (`LOGIN` / `AUTHENTICATE PLAIN`) не прошёл + (NO LOGIN failed). himalaya логинится успешно — вероятно, проблема в + способе авторизации/лимитах для «незнакомого» клиента, а не в пароле. + В design — использовать существующий himalaya-конфиг и проверенный путь. +- Текущий архиватор: `scripts/mail_archive.py` — poll каждые 5 мин (Hermes cron + `mail-archive-every-5min`, id 5f2305b2bbf8). Отслеживает `last_uid` по папкам + в `/opt/hermes/email/state/mail-archive-last-.json`. + Уже есть: `fetch_attachments_imaplib()` (сырой IMAP, BODY.PEEK[], не ставит + `\Seen`), `get_inbox_subfolders()` (динамическое обнаружение 136 папок, + CHILDREN), `_imap_utf7_encode()` (modified UTF-7 для кириллических папок). +- Архив: `/opt/hermes/email//YYYY/MM//email.md` (frontmatter: + id, folder, subject, from, to, date, flags, message_id, in_reply_to, references, + cc, content_type; + classification, handled_*, has_attachment). + **5467 писем** на 2026-09-15. +- Индексы: `scripts/mail_index.py` (SQLite FTS5 `/opt/hermes/email/mail_index.db`), + `sqlite_search.py`, классификатор `email_classifier.py` (Qwen3:8b, Ollama), + обработчики `email_handlers.py` (urgent→TG, task→VTODO, meeting→VEVENT), + `digest.py` (еженедельный). +- Пользователь активно меняет ящик в клиенте: переносит в папки, удаляет, + отвечает. Эти изменения сейчас НЕ видны архиватору. + +## Решение + +### 1. Выбор модели: 2 постоянных процесса (systemd user units) + +Не один процесс на всё, а два — разделение ответственности (как ты предложил): + +1. **`email-imap-stream.service`** — держит IMAP-соединение, IDLE, пишет + события в SQLite (`mailbox_events`), архивирует новые письма. + *Единственный* процесс, который разговаривает с IMAP (кроме fallback-cron). +2. **`email-change-analyzer.service`** — постоянно читает `mailbox_events`, + обновляет online-копию (`mailbox_state`), запускает классификатор/обработчики, + индексы, может писать «журнал изменений» для пользователя. + +Почему 2, а не 1: падение анализатора не теряет события (они в ChangeLog, +append-only); падение stream-процесса не останавливает анализ (fallback-cron +докачает письма). Отвязка скорости IMAP от скорости LLM-анализа. + +### 2. Библиотека для IMAP stream + +Варианты: +- (a) `aioimaplib` — Python asyncio IMAP с поддержкой IDLE, reconnect, ивенты. + Мягкая зависимость; если нет — pip install. Совместим с Exchange (IDLE есть). +- (b) Сырой socket (как уже сделано в `fetch_attachments_imaplib`) + select() + на IDLE — без новых зависимостей, но больше кода (reconnect, парсинг + untagged-ответов, литералы). + +**Выбор: (a) `aioimaplib`** — проверенная библиотека, IDLE/reconnect из коробки; +сырой IMAP оставляем только в `fetch_attachments_imaplib` (там он уже работает +и менять не нужно). Если `aioimaplib` недоступен/не заводится на нашем Python — +fallback (b). + +**Авторизация:** переиспользовать учётку из `config/himalaya-config.toml` +(login `e.storozhenko`, host `mail.corpoffice.tech`). Для stream-процесса — +`AUTH=PLAIN` с SASL-IR (Exchange поддерживает; но в тесте не прошло — проверить +в задаче 2: использовать himalaya `account sync`/`envelope list` как эталон, +возможно, потребуется `\r\n`-формат или конкретный порядок SASL). + +### 3. Схема SQLite (новая БД или таблицы в существующей) + +Отдельная БД: `/opt/hermes/email/state/mailbox.db` (не смешивать с +`mail_index.db`, чтобы не ломать FTS-индексы и было легко откатиться). + +```sql +-- Онлайновая копия ящика: текущее состояние каждого письма +CREATE TABLE mailbox_state ( + uid INTEGER NOT NULL, -- UID в папке + folder TEXT NOT NULL, -- папка (например 'INBOX', 'INBOX/Проекты') + message_id TEXT, -- Message-ID + in_reply_to TEXT, + references TEXT, -- цепочка (References / In-Reply-To) + subject TEXT, + from_addr TEXT, + date TEXT, + flags TEXT, -- JSON-массив флагов IMAP + has_attachment INTEGER DEFAULT 0, + archive_path TEXT, -- путь к email.md (если заархивировано) + etag TEXT, -- ETag сервера (возможно, не нужен для IMAP) + last_seen TEXT, -- ISO-время последней синхронизации + deleted INTEGER DEFAULT 0, -- soft-delete (пользователь удалил на сервере) + PRIMARY KEY (folder, uid) +); + +-- Журнал изменений (ChangeLog), append-only +CREATE TABLE mailbox_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ts TEXT NOT NULL, -- ISO-время события + event TEXT NOT NULL, -- added | moved | deleted | flag_changed | replied | reconcile + folder TEXT, + uid INTEGER, + message_id TEXT, + details TEXT -- JSON: from_folder, to_folder, old_flags, new_flags и т.п. +); +CREATE INDEX idx_events_message ON mailbox_events(message_id); +CREATE INDEX idx_events_ts ON mailbox_events(ts); +``` + +### 4. Алгоритм stream-процесса (`imap_stream.py`) + +``` +loop: + 1. connect + login (AUTH=PLAIN / STARTTLS) + 2. folder list (get_inbox_subfolders через IMAP, dynamic) + 3. для каждой папки: SELECT; если нет в mailbox_state — полный reconcile + (UID FETCH 1:* (FLAGS UID MESSAGE-ID ...) = initial sync) + 4. reconcile источником истины: для INBOX + подпапок + - UID FETCH (новые UID > last_uid) → событие added, архивация + - сравнение FLAGS (в т.ч. \Seen, \Answered, \Flagged) → flag_changed + - UID SEARCH EXPUNGE / отсутствие в SELECT → deleted (если был в state) + 5. IDLE loop (для INBOX, и по очереди для подпапок если CPU позволяет): + - IDLE → ждём untagged: EXISTS (новое письмо), EXPUNGE (удаление), + FETCH FLAGS (изменение флагов), MOVE (если сервер шлёт) + - на событие: обработать (архивировать/обновить state/записать event) + - продлевать IDLE каждые ~29 мин (сервер обычно рвёт после 30) + 6. при обрыве: reconnect + reconcile (REQ-IMAP-SYNC-001/002) + 7. период паузы между reconcile: 60-300 сек (конфигурируемо) +``` + +**Архивация новых писем** — переиспользуем существующий код `mail_archive.py`: +`get_email_content()` (himalaya `message read --preview`, не ставит Seen) + +`fetch_attachments_imaplib()` (BODY.PEEK[]). Функции вынести в общий модуль +или импортировать (mail_archive.py уже модульный). + +**Удаление:** при событии EXPUNGE/`deleted` — НЕ удаляем файл email.md (письмо +остаётся в архиве, REQ-IMAP-SYNC-004 soft-delete), помечаем `deleted=1` в +`mailbox_state`, пишем событие `deleted` в ChangeLog. Пользователь сможет видеть +«это письмо вы удалили» в RAG/поиске. + +**Перемещение:** при обнаружении (MOVE/UIDPLUS от сервера или reconcile: +UID в старой папке исчез, в новой появился с тем же Message-ID) — обновить +`folder`, записать `moved` с from/to. Если UID меняется (Exchange MOVE) — +следить по Message-ID. + +### 5. Анализатор (`change_analyzer.py`) + +``` +loop: + 1. читать mailbox_events от последнего обработанного id (offset в state) + 2. для каждого события: + added → классифицировать (email_classifier), обработать (email_handlers), + обновить индексы (mail_index --incremental) + flag_changed → обновить flags в mailbox_state; если \Answered → событие replied + moved → обновить folder; если классификация была — можно переклассифицировать + deleted → пометить в RAG/индексе удалённым (не удаляя файл) + 3. записать обработанный id (persistent) + 4. спать 1-5 сек (или ждать сигнала от stream через очередь) +``` + +**Связь stream ↔ analyzer:** через SQLite (ChangeLog) — это проще и надёжнее, +чем IPC/очереди. Stream пишет, analyzer читает с offset. Несколько analyser +процессов не нужны (одна учётка). + +### 6. Управление / systemd + +Два user unit (в `~/.config/systemd/user/`): + +```ini +# email-imap-stream.service +[Unit] +Description=Email IMAP realtime sync stream +After=network-online.target + +[Service] +Type=simple +ExecStart=/usr/bin/python3 /opt/hermes/email-assistant/scripts/imap_stream.py +Restart=always +RestartSec=5 +Environment=HOME=/home/estorozhenko # для himalaya (конфиг там) + +[Install] +WantedBy=default.target +``` + +Аналогично `email-change-analyzer.service`. Включить: `systemctl --user enable --now`. + +**Hermes cron** `mail-archive-every-5min` → оставить как **fallback**: +если оба сервиса не работают (например, после ребута без автозапуска), cron +всё равно архивирует раз в 5 мин (REQ-IMAP-SYNC-007). Чтобы не дублировать: +stream-процесс и cron оба идемпотентны (last_uid / mailbox_state). + +### 7. Что НЕ делаем (ограничения) + +- **НЕ** переписываем архиватор под IMAP с нуля — используем существующий + `mail_archive.py` как библиотеку/подпроцесс. +- **НЕ** удаляем файлы писем при удалении на сервере (RAG должен видеть + «удалено», но данные не теряем). +- **НЕ** реалтайм-синхронизация «один к одному» всех 137 папок через отдельные + IDLE-коннекты — IDLE держим на INBOX (главный источник), подпапки — через + reconcile (раз в 1-5 мин) + CHILDREN-обнаружение. Это баланс скорости и + ресурсов (Exchange лимитирует коннекты на учётку). +- **НЕ** реализуем SMTP/отправку — только чтение (как сейчас). + +### 8. Риски и смягчение + +| Риск | Смягчение | +|---|---| +| Exchange рвёт IDLE после ~30 мин | авто-reconnect + reconcile после каждого обрыва | +| LOGIN с чистого скрипта не прошёл | использовать himalaya-путь; AUTH=PLAIN; задача 2 — проверить формат | +| IDLE не гарантирует доставку всех событий | reconcile каждые 1-5 мин (REQ-002) | +| MOVE меняет UID на Exchange | следить по Message-ID, а не UID | +| Много папок (137) — много коннектов | IDLE только на INBOX; подпапки — reconcile по очереди | +| aioimaplib не установлен | fallback на сырой socket (уже есть паттерн) | \ No newline at end of file diff --git a/openspec/changes/imap-realtime-sync/proposal.md b/openspec/changes/imap-realtime-sync/proposal.md new file mode 100644 index 0000000..693abc5 --- /dev/null +++ b/openspec/changes/imap-realtime-sync/proposal.md @@ -0,0 +1,150 @@ +# Proposal: Потоковая синхронизация почты (IMAP IDLE) и онлайновая копия ящика + +## Why + +Сейчас `mail_archive.py` опрашивает IMAP **каждые 5 минут** по расписанию (Hermes cron). +Это pull-модель с двумя фундаментальными ограничениями: + +1. **Задержка до 5 минут** — новые письма попадают в архив и в RAG-индексы не сразу, + а в лучшем случае через 5 минут (а с учётом очереди классификатора/обработчиков — дольше). +2. **Слепота к изменениям, которые делает сам пользователь.** Пользователь активно + работает с почтой в клиенте: переносит письма между папками, удаляет, отвечает, + помечает прочитанными. Архиватор знает только `last_uid` и **не видит**: + - перемещение письма из INBOX в подпапку (письмо «пропадает» из INBOX, но в архиве остаётся в INBOX) + - удаление письма (архив хранит удалённое письмо, но не знает, что оно удалено) + - ответ/флаг `\Answered`, `\Flagged` (архив хранит статичные флаги) + - появление новых писем в **существующих** папках (last_uid инкрементален, но только + для писем, добавленных после последнего опроса) + +Из-за этого **нельзя строить RAG/граф знаний по почте корректно**: цепочки +«письмо пришло → прочитано → перемещено → на него ответили» в данных отсутствуют, +и при поиске агент не может сказать «это письмо вы удалили» — он просто не находит +его в той папке, где оно было при архивации. + +**Хочется:** постоянное (потоковое) соединение с IMAP — как в почтовом клиенте — +которое **мгновенно** узнаёт о новых письмах (через IDLE-уведомления), и отдельный +процесс, который ведёт **онлайновую копию ящика** с полной историей изменений +(лог событий: письмо добавлено/перемещено/удалено/прочитано/получен ответ). +На этой основе потом можно: корректно строить RAG, писать «журнал изменений», +отвечать на вопросы про историю ящика. + +## What Changes + +### Архитектура: 2 сервиса (микросервисы на одной машине) + +Заменяем «cron-скрипт каждые 5 минут» на два **постоянных процесса** (systemd units): + +``` +┌─────────────────────────────┐ ┌─────────────────────────────┐ +│ svc 1: IMAP Stream/Sync │ │ svc 2: Change Analyzer │ +│ (постоянный IDLE-коннект) │ │ (анализ изменений) │ +│ │ │ │ +│ • держит 1+N IDLE-соединений│ │ • читает журнал изменений │ +│ • мгновенно видит события │ │ • обновляет online-копию │ +│ • пишет в ChangeLog │ │ • RAG/классификация/ │ +│ • скачивает новые письма │ │ обработчики │ +└───────────┬─────────────────┘ └─────────────┬───────────────┘ + │ журнал изменений (append-only) │ + ▼ ▼ + ┌──────────────────────────────────────────────────────┐ + │ SQLite: online-копия ящика (состояние) │ + │ + ChangeLog (события, append-only) │ + └──────────────────────────────────────────────────────┘ +``` + +**svc 1 — IMAP Stream Sync («синхронизатор»):** +- одно постоянное соединение с IMAP (Exchange, :143 STARTTLS), авторизация как у himalaya; +- `IDLE`-команда на INBOX (и ключевых папках) — сервер **пушит** уведомления + `* N EXISTS` / `* N EXPUNGE` / `* N FETCH FLAGS` сразу при изменении; +- на событие — архив письма (через существующий `fetch_attachments_imaplib`, BODY.PEEK[], + не ставя `\Seen`), обновление online-копии, запись события в ChangeLog; +- периодический (раз в N мин) **reconcile** — полный `UID FETCH` изменений с сервера + (страховка от пропущенных событий при обрыве IDLE: сервер не гарантирует доставку всех + событий через IDLE, но reconcile это закрывает); +- папки: динамическое обнаружение (`get_inbox_subfolders()` уже есть) + CHILDREN; +- **MOVE/UIDPLUS** (поддерживается Exchange) — позволяет точнее отслеживать перемещения + (`UID MOVE` возвращает старый/новый UID). + +**svc 2 — Change Analyzer («анализатор»):** +- постоянно крутится, читает ChangeLog из SQLite (или файловый след); +- поддерживает **online-копию ящика**: для каждого письма — текущий folder, flags, uid, + thread (References/In-Reply-To), время последнего изменения; +- корректно строит **цепочки**: письмо пришло → прочитано (flag) → перемещено в папку → + на него ответили (по References) → удалено (если пользователь удалил); +- на основе изменений обновляет RAG/индексы и запускает классификатор/обработчики + (аналог текущего `mail-classify-handlers` cron, но по событиям, а не по расписанию); +- **журнал изменений** = история «что, когда, с каким письмом произошло» — его можно + показывать пользователю и использовать в RAG (в отличие от текущей модели, где + у нас только статичные снимки). + +### Новые компоненты + +- `scripts/imap_stream.py` — svc 1 (постоянный процесс; systemd unit `email-imap-stream.service`) +- `scripts/change_analyzer.py` — svc 2 (анализ ChangeLog → online-копия → RAG/обработчики) +- `schema`: таблица `mailbox_state` (online-копия) + `mailbox_events` (ChangeLog) в + существующей SQLite (`/opt/hermes/email/mail_index.db`) или отдельной + `/opt/hermes/email/state/mailbox.db` +- systemd units (user-level) вместо Hermes cron; cron остаётся как fallback/страховка + (например, каждые 5 минут — reconcile, если оба сервиса умерли) + +### Существующие скрипты не ломаются + +- `mail_archive.py` остаётся (он уже умеет BODY.PEEK[] и не ставит Seen); + svc 1 может использовать его функции как библиотеку (или дублировать минимально). +- `email_classifier.py`, `email_handlers.py`, `mail_index.py`, `digest.py` — вызываются + из svc 2 (по событиям) и/или остаются по cron (по расписанию) — поведение не меняется. + +## Capabilities + +### New Capabilities + +- `imap-realtime-sync`: Постоянное IMAP-соединение (IDLE) с мгновенным обнаружением + новых писем и изменений (перемещение/удаление/флаги) — без опроса по расписанию. +- `mailbox-online-copy`: Онлайновая копия ящика (folder/flags/uid/thread у каждого письма) + + журнал изменений (ChangeLog: добавить/переместить/удалить/прочитать/ответить) — + источник истины для RAG и ответов «что случилось с письмом». +- `email-thread-tracking`: Отслеживание цепочек писем (References/In-Reply-To) и событий + жизни письма (пришло → прочитано → перемещено → ответ → удалено). + +### Modified Capabilities + +- `email-attachments`: скачивание вложений теперь инициируется событиями IDLE, + а не только опросом по расписанию (механика та же: BODY.PEEK[]). +- `email-classification` / `email-handlers`: запускаются по событиям ChangeLog + (новое письмо/изменение), а не только по cron. + +## Impact + +- **Процессы:** 2 новых systemd user unit (`email-imap-stream.service`, + `email-change-analyzer.service`), всегда запущены. Hermes cron `mail-archive-every-5min` + → заменяется на reconcile-cron (или остаётся как fallback). +- **Сеть:** постоянное TCP-соединение с mail.corpoffice.tech:143 (keep-alive, IDLE + продлевается каждые ~29 мин). Раньше соединение открывалось каждые 5 минут. +- **Данные:** новая таблица SQLite (online-копия + ChangeLog). Письма в `/opt/hermes/email/` + не переезжают — формат `email.md` и структура папок сохраняются (совместимость с + `mail_index.py`, Obsidian-бэкапами, поиском). +- **Риски:** + - Exchange может рвать IDLE-соединения (keep-alive не вечен) — нужен auto-reconnect + с reconcile после переподключения (обязательное требование). + - IDLE на Exchange поддерживается (проверено: в CAPABILITY есть `IDLE`), но поведение + сервера при `EXPUNGE`/перемещении может отличаться — reconcile закрывает. + - Логин: himalaya логинится успешно. Разобрано 2026-09-15: отказ LOGIN + с чистого скрипта — это rate-limit Exchange (после серии быстрых попыток сервер + молчит timeout), а не TLS-fingerprinting. Работает простой `LOGIN` из + `imap_connect()` с re-try (экспоненциальная пауза). Подробности — в design.md. + - Два постоянных процесса = чуть больше памяти/CPU, но для одной учётки это копейки. +- **Наблюдаемость:** авторизация и метрики (доступность сервера, успешность, + активная сессия) логируются в `/opt/hermes/email/logs/imap_client.log` + (`imap_log()`/`imap_metrics()`/`imap_session` в `scripts/imap_client.py`). + Полноценный Prometheus/статус-эндпоинт для stream-сервиса — позднее (tasks.md). +- **Документация:** README/STATUS/WALKTHROUGH — новые сервисы, порты (нет новых внешних), + как перезапускать, как смотреть журнал изменений. + +## Rollback + +1. Остановить systemd units (`systemctl --user stop email-imap-stream email-change-analyzer`). +2. Вернуть Hermes cron `mail-archive-every-5min` (он никуда не делся — просто выключен). +3. Удалить новую таблицу/базу (online-копия) — она производная, письма в `/opt/hermes/email/` + не затрагиваются. +4. Существующие скрипты (`mail_archive.py`, классификатор, обработчики) не меняют поведение + при отказе от сервисов — они работают как раньше. \ No newline at end of file diff --git a/openspec/changes/imap-realtime-sync/specs/email-attachments/spec.md b/openspec/changes/imap-realtime-sync/specs/email-attachments/spec.md new file mode 100644 index 0000000..2ac262d --- /dev/null +++ b/openspec/changes/imap-realtime-sync/specs/email-attachments/spec.md @@ -0,0 +1,27 @@ +# email-attachments Specification + +## MODIFIED Requirements + +### Requirement: Вложения сохраняются в каталог письма + +Для каждого письма с вложениями (флаг `has_attachment: true` в frontmatter) +вложения MUST быть сохранены в подкаталог `attachments/` каталога письма +(`/opt/hermes/email//YYYY/MM//attachments/`). Скачивание выполняется +как при опросе по расписанию (текущее поведение), так и при событийной +синхронизации (сервис `imap-realtime-sync`). Скачивание вложений MUST NOT +выставлять флаг `\Seen`. + +#### Scenario: Письмо с вложением архивировано +- **WHEN** `mail_archive.py` заархивировал письмо с `has_attachment: true` +- **THEN** файлы вложений лежат в `/attachments/` и совпадают с вложениями на IMAP-сервере + +#### Scenario: Вложения при событийной синхронизации +- **GIVEN** новое письмо с вложением обнаружено сервисом `email-imap-stream` +- **WHEN** письмо архивируется по событию IDLE +- **THEN** вложения скачаны в `/attachments/`, флаг `\Seen` не выставлен + +#### Scenario: Вложения при reconcile +- **GIVEN** письмо с вложением было пропущено (обрыв соединения) +- **WHEN** сервис выполняет reconcile +- **THEN** вложения докачиваются в `/attachments/` (то же поведение, + что существующий `--attachments-backfill`) \ No newline at end of file diff --git a/openspec/changes/imap-realtime-sync/specs/email-classification/spec.md b/openspec/changes/imap-realtime-sync/specs/email-classification/spec.md new file mode 100644 index 0000000..d756ba1 --- /dev/null +++ b/openspec/changes/imap-realtime-sync/specs/email-classification/spec.md @@ -0,0 +1,26 @@ +# email-classification Specification + +## MODIFIED Requirements + +### Requirement: Классификация каждого нового письма + +Классификация писем (Qwen3:8b через Ollama) MUST запускаться не только по +расписанию (cron `mail-classify-handlers`), но и по событиям из ChangeLog +(новое письмо / изменённое письмо), которые генерирует сервис +`email-imap-realtime-sync`. Механика классификации (теги info/urgent/task/meeting ++ `classification`/`classification_reason` в frontmatter) не меняется. + +#### Scenario: Новое письмо после архивации +- **GIVEN** письмо заархивировано (`mail_archive.py`) +- **WHEN** `email_classifier.py` обрабатывает письмо +- **THEN** в frontmatter появляются `classification` и `classification_reason` + +#### Scenario: Новое письмо классифицируется сразу +- **GIVEN** в ChangeLog появилось событие `added` +- **WHEN** анализатор `email-change-analyzer` обрабатывает событие +- **THEN** письмо классифицируется (тег + обоснование в frontmatter) без ожидания cron + +#### Scenario: Событие обработано повторно (идемпотентность) +- **GIVEN** письмо уже классифицировано (есть `classification` в frontmatter) +- **WHEN** анализатор видит то же событие повторно (например, после рестарта) +- **THEN** классификация не выполняется заново (идемпотентно, как сейчас) \ No newline at end of file diff --git a/openspec/changes/imap-realtime-sync/specs/imap-realtime-sync/spec.md b/openspec/changes/imap-realtime-sync/specs/imap-realtime-sync/spec.md new file mode 100644 index 0000000..c9af513 --- /dev/null +++ b/openspec/changes/imap-realtime-sync/specs/imap-realtime-sync/spec.md @@ -0,0 +1,118 @@ +# imap-realtime-sync Specification + +## ADDED Requirements + +### Requirement: REQ-IMAP-SYNC-001: Постоянное IMAP-соединение с IDLE + +Система MUST поддерживать постоянное IMAP-соединение с почтовым сервером +(mail.corpoffice.tech:143, STARTTLS) с использованием команды `IDLE` +(поддерживается сервером, проверено в CAPABILITY), чтобы узнавать о новых +письмах и изменениях в ящике **мгновенно**, без опроса по расписанию. + +#### Scenario: Новое письмо появляется в INBOX +- **GIVEN** сервис `email-imap-stream` запущен и держит IDLE-соединение с INBOX +- **WHEN** в INBOX приходит новое письмо +- **THEN** сервис получает уведомление `* N EXISTS` от сервера в течение секунд + (не дольше таймаута IDLE) И инициирует архивацию письма + +#### Scenario: IDLE-соединение обрывается +- **GIVEN** IDLE-соединение активно +- **WHEN** сервер закрывает соединение (таймаут/сбой) +- **THEN** сервис автоматически переподключается и выполняет полный reconcile + (синхронизацию изменений), чтобы не пропустить события, случившиеся во время обрыва + +### Requirement: REQ-IMAP-SYNC-002: Reconcile как страховка от пропущенных событий + +Сервис MUST периодически (не реже 1 раза в 5 минут) выполнять reconcile — +полную сверку состояния ящика с сервером (UID FETCH / STATUS), потому что IDLE +не гарантирует доставку всех событий (особенно при длительных соединениях). + +#### Scenario: Событие потеряно при обрыве +- **GIVEN** сервис работал, но IDLE оборвался на 3 минуты, за это время письмо + было перемещено пользователем +- **WHEN** сервис переподключается и делает reconcile +- **THEN** перемещение обнаруживается и фиксируется в ChangeLog + +### Requirement: REQ-IMAP-SYNC-003: Архивация без установки \Seen + +Любое скачивание тела/вложений (по событию IDLE или reconcile) MUST NOT +выставлять IMAP-флаг `\Seen`. Используется существующий +`fetch_attachments_imaplib()` (BODY.PEEK[]) и чтение с `--preview`. + +#### Scenario: Новое письмо архивируется по событию +- **GIVEN** новое письмо в INBOX с флагами `()` (непрочитанное) +- **WHEN** сервис по событию IDLE архивирует письмо +- **THEN** на IMAP флаги письма остаются `()` (письмо остаётся непрочитанным) + +### Requirement: REQ-IMAP-SYNC-004: Онлайновая копия ящика + +Система MUST вести онлайновую копию ящика в SQLite (`mailbox_state`): для +каждого письма — текущий folder, UID, флаги, Message-ID, References/In-Reply-To, +дата и время последнего изменения. Копия обновляется при каждом событии. + +#### Scenario: Письмо перемещено пользователем +- **GIVEN** письмо было в INBOX, пользователь переместил его в `INBOX/Проекты` +- **WHEN** сервис получает событие пересмещения (MOVE/UIDPLUS или reconcile) +- **THEN** в `mailbox_state` folder обновлён на `INBOX/Проекты`, в ChangeLog + записано событие `moved` + +#### Scenario: Письмо удалено пользователем +- **GIVEN** письмо было в ящике +- **WHEN** пользователь удаляет письмо (EXPUNGE/FLAGS \Deleted) +- **THEN** в `mailbox_state` письмо помечается удалённым (soft-delete), в ChangeLog + записано событие `deleted`; файл письма в архиве сохраняется + +### Requirement: REQ-IMAP-SYNC-005: Журнал изменений (ChangeLog) + +Система MUST вести append-only журнал изменений (`mailbox_events`): каждое +событие (added / moved / deleted / flag_changed / replied) с timestamp, folder, +UID, Message-ID. Журнал служит источником истины для анализатора и RAG. + +#### Scenario: Запись о новом письме +- **GIVEN** сервис обнаружил новое письмо +- **WHEN** письмо заархивировано +- **THEN** в `mailbox_events` создана запись `added` с полным контекстом (UID, + folder, Message-ID, timestamp) + +#### Scenario: Запись об ответе +- **GIVEN** пользователь ответил на письмо (флаг \Answered или новое письмо с In-Reply-To) +- **WHEN** сервис видит изменение +- **THEN** в ChangeLog записано событие `replied` (по References/In-Reply-To) или + `flag_changed` (для \Answered) + +### Requirement: REQ-IMAP-SYNC-006: Анализатор изменений + +Система MUST иметь отдельный процесс (`email-change-analyzer`), который читает +ChangeLog и обновляет производные данные: online-копию, RAG-индексы, +классификацию/обработчики. Анализатор работает по событиям (потоково), а не +по расписанию. + +#### Scenario: Новое письмо → классификация и обработчики +- **GIVEN** в ChangeLog появилось событие `added` +- **WHEN** анализатор обрабатывает событие +- **THEN** запускается классификация (Qwen3:8b) и обработчики (urgent→TG, + task→VTODO, meeting→VEVENT) — как сейчас, но по событию, без ожидания cron + +#### Scenario: Письмо удалено → RAG-индексы +- **GIVEN** в ChangeLog событие `deleted` +- **WHEN** анализатор обрабатывает событие +- **THEN** RAG/индекс помечает письмо удалённым (не удаляя файл) — при поиске + агент может сказать «это письмо удалено пользователем» + +### Requirement: REQ-IMAP-SYNC-007: Управление как systemd-сервисами + +Сервисы MUST работать как постоянные процессы под systemd (user units), +автозапуск при логине, auto-restart при падении, логи в journald. Hermes cron +остаётся как fallback-reconcile (страховка, если оба сервиса не работают). + +#### Scenario: Сервис упал +- **GIVEN** `email-imap-stream.service` запущен +- **WHEN** процесс падает/убивается +- **THEN** systemd перезапускает его автоматически (Restart=always), при старте — + reconcile + +#### Scenario: Оба сервиса не работают +- **GIVEN** оба сервиса остановлены (например, после перезагрузки без автозапуска) +- **WHEN** проходит 5 минут +- **THEN** Hermes cron (fallback) запускает `mail_archive.py` — архивация не + останавливается полностью \ No newline at end of file diff --git a/openspec/changes/imap-realtime-sync/tasks.md b/openspec/changes/imap-realtime-sync/tasks.md new file mode 100644 index 0000000..70df3e7 --- /dev/null +++ b/openspec/changes/imap-realtime-sync/tasks.md @@ -0,0 +1,78 @@ +# Tasks: Потоковая синхронизация почты (IMAP IDLE) и онлайновая копия ящика + +## Задачи + +### Фаза 0: Подготовка и проверка авторизации +- [ ] 1. Установить/проверить `aioimaplib` (pip) — библиотека для IDLE-потока; + если не ставится — зафиксировать fallback на сырой socket. + (2026-09-15: `aioimaplib` НЕ установлен; на следующем шаге — решить + ставить или идти на сыром socket, т.к. imap_client.py уже это умеет.) +- [x] 2. **Проверить авторизацию** для stream-процесса: + himalaya логинится успешно (эталон), чистый IMAP-скрипт — нет. + Отработать AUTH=PLAIN (SASL-IR), формат login; зафиксировать рабочий + вариант в коде. (Блокер — без него stream не запустится.) + **Выполнено 2026-09-15**: создан `scripts/imap_client.py` с отдельной + функцией `imap_connect()` (STARTTLS + LOGIN + re-try против rate-limit). + Реальная проверка: LOGIN как e.storozhenko прошёл, `fetch_attachments_imaplib` + через неё скачал `Переместить стол.docx` 1.1 МБ (UID 14200) без \Seen. + Rate-limit Exchange: после серии быстрых попыток сервер молчит (timeout), + поэтому в `imap_connect()` re-try с экспоненциальной паузой. +- [x] 2b. **Логирование авторизации + метрики доступности/сессии**: + `imap_log()` пишет JSON-строки в `/opt/hermes/email/logs/imap_client.log` + (conn_ok / conn_error / auth_ok / auth_failed / session_started / session_ended); + `imap_metrics()` считает auth_success_rate, conn_error, sessions_active; + `imap_session` — контекстный менеджер (гарантирует session_ended). + Проверено: `python3 imap_client.py --metrics` показывает метрики. + Пароль никогда не логируется. + **TODO (позже)**: полноценный мониторинг — Prometheus-формат/статус-эндпоинт + для stream-сервиса, алерты при падении auth_success_rate / conn_error. + +### Фаза 1: Скелет сервиса imap_stream.py +- [ ] 3. `scripts/imap_stream.py`: connect + STARTTLS + login; folder list + (get_inbox_subfolders); SELECT INBOX; IDLE-цикл с обработкой untagged + (EXISTS/EXPUNGE/FETCH FLAGS); reconnect при обрыве. +- [ ] 4. Реализовать reconcile: UID FETCH новых писем (>last_uid) → событие + `added` + архивация (переиспользовать mail_archive.py); сравнение FLAGS → + `flag_changed`; отсутствие после EXPUNGE → `deleted` (soft). +- [ ] 5. SQLite schema: `mailbox_state` + `mailbox_events` (ChangeLog) — + создать `/opt/hermes/email/state/mailbox.db`, функции init/insert/read. + +### Фаза 2: Change Analyzer +- [ ] 6. `scripts/change_analyzer.py`: чтение mailbox_events от offset; + на `added` — классификация (email_classifier) + обработчики (email_handlers); + на `moved`/`deleted`/`flag_changed` — обновление mailbox_state, RAG-метка. + **Уведомления в ЛС Telegram уже проверены живьём (2026-09-15)**: письмо 3216 + помечено urgent → доставка в private chat 281328953 (kpa39l), msg_id=60, + подтверждено пользователем. Для ЛС нужен явный `TELEGRAM_CHAT_ID=281328953` + (дефолт в email_handlers.py — канал @dedinit_vesti). +- [ ] 7. Обработка `replied`: \Answered или новое письмо с In-Reply-To на + известный Message-ID → событие `replied`; обновить thread в state. + +### Фаза 3: systemd и интеграция +- [ ] 8. systemd user units (`email-imap-stream.service`, + `email-change-analyzer.service`), автозапуск; проверка auto-restart. +- [ ] 9. Hermes cron `mail-archive-every-5min` → fallback (идемпотентен с + stream; не дублирует). Проверить, что при работающем stream cron не + архивирует повторно (last_uid / state). +- [ ] 10. Живой тест: новое письмо (отправить себе/ждущее), перемещение + (через IMAP MOVE/клиент), удаление, ответ — всё фиксируется в ChangeLog + и mailbox_state; \Seen не ставится (проверка флагов до/после). + +### Фаза 4: Документация и доводка +- [ ] 11. STATUS.md / README / WALKTHROUGH: сервисы, порты (нет новых внешних), + как смотреть ChangeLog (`sqlite3 mailbox.db 'select * from mailbox_events'`), + как перезапускать. +- [ ] 12. RAG-интеграция (опционально, Фаза 2 проекта): индексация событий + ChangeLog в Qdrant — чтобы агент мог ответить «что случилось с письмом». + +## Верификация + +- `systemctl --user status email-imap-stream email-change-analyzer` — active (running) +- `sqlite3 /opt/hermes/email/state/mailbox.db 'select count(*) from mailbox_state'` — растёт +- `sqlite3 ... 'select event, count(*) from mailbox_events group by event'` — + есть added/moved/deleted/flag_changed +- Новое письмо в INBOX → архив появляется в течение ~1-2 мин (не 5) +- Перемещение письма в клиенте → в mailbox_state folder обновлён, в ChangeLog `moved` +- Удаление письма → `deleted` в ChangeLog, файл email.md на месте (soft-delete) +- Ответ → `replied` или `flag_changed` (\Answered) +- Флаги непрочитанного письма на IMAP после архивации: `()` → `()` (Seen нет) \ No newline at end of file diff --git a/scripts/email_handlers.py b/scripts/email_handlers.py index af9b5e3..eb611fe 100644 --- a/scripts/email_handlers.py +++ b/scripts/email_handlers.py @@ -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" diff --git a/scripts/imap_client.py b/scripts/imap_client.py new file mode 100644 index 0000000..57051c8 --- /dev/null +++ b/scripts/imap_client.py @@ -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) \ No newline at end of file diff --git a/scripts/mail-classify-handlers.sh b/scripts/mail-classify-handlers.sh index 325d414..3bffb2b 100755 --- a/scripts/mail-classify-handlers.sh +++ b/scripts/mail-classify-handlers.sh @@ -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 \ No newline at end of file +# Классификатор: до 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 \ No newline at end of file diff --git a/scripts/mail_archive.py b/scripts/mail_archive.py index c18e05d..717e555 100755 --- a/scripts/mail_archive.py +++ b/scripts/mail_archive.py @@ -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.