mirror of
https://gitverse.ru/kpa39l/email-assistant.git
synced 2026-09-29 09:15:09 +00:00
Compare commits
7 Commits
3e7b1dd618
..
master
| Author | SHA1 | Date | |
|---|---|---|---|
| b34993402f | |||
| 83bbe3277d | |||
| 99e8e847f3 | |||
| 5a81c86599 | |||
| defcf53d75 | |||
| 65c167acc2 | |||
| 0d86f843c3 |
@@ -6,6 +6,8 @@
|
||||
Все коммиты и изменения проекта пушатся именно туда.
|
||||
- **Gitea (gitea.nixg.ru) — ЗЕРКАЛО (mirror).** Не источник; используется
|
||||
как вторичная/внутренняя копия. Не пушить в gitea как в основное хранилище.
|
||||
Pull mirror: `estorozhenko/email-assistant` ← gitverse `kpa39l/email-assistant`,
|
||||
интервал 12h (создан 2026-09-15).
|
||||
|
||||
## URL-ы
|
||||
|
||||
@@ -19,6 +21,30 @@
|
||||
IdentityFile ~/.ssh/gitverse, IdentitiesOnly yes) — валиден для аутентификации,
|
||||
но push-to-create отключён; для создания репо — API.
|
||||
|
||||
## Токены (местоположение)
|
||||
|
||||
**Единый источник — `/opt/hermes/.hermes/secrets/git-tokens.env`** (chmod 600, единственный
|
||||
истинный файл токенов; НЕ коммитится ни в один репозиторий):
|
||||
|
||||
| Переменная | Назначение |
|
||||
|---|---|
|
||||
| `GITVERSE_LOGIN` | логин gitverse (`kpa39l`, НЕ estorozhenko) |
|
||||
| `GITVERSE_API` | база gitverse API (`https://api.gitverse.ru`) |
|
||||
| `GITVERSE_PAT` | PAT gitverse (создан 19 июл 2026, имя `hermes-agent`) |
|
||||
| `GITEA_API` | база gitea API (`http://127.0.0.1:3000/api/v1`) |
|
||||
| `GITEA_TOKEN` | токен gitea (читает/пишет репо, зеркала) |
|
||||
|
||||
Загрузка: `set -a && source /opt/hermes/.hermes/secrets/git-tokens.env && set +a`
|
||||
|
||||
Прочие токены проекта (НЕ в git-tokens.env):
|
||||
- **Telegram Bot (уведомления о важных письмах):** `/opt/vesti/.env` → `VESTI_BOT_TOKEN`
|
||||
(защищён Hermes-секретом — не читать/не обходить; использовать через `source` в рантайме).
|
||||
- **TELEGRAM_CHAT_ID** (ЛС kpa39l, `281328953`): `/opt/hermes/email-assistant/.env` (gitignored).
|
||||
- **SSH-ключ gitverse:** `~/.ssh/gitverse` (алиас `Host gitverse.ru` в `~/.ssh/config`).
|
||||
|
||||
⚠️ Токены НЕ хранить в коде/README/коммитах — только в secrets-файле и remote URL git.
|
||||
Если в репозитории случайно закоммичен токен — немедленно отозвать и перевыпустить.
|
||||
|
||||
## Как пушить
|
||||
|
||||
```bash
|
||||
@@ -58,4 +84,5 @@ HTTP 201 = создан; 409 = уже существует. `auto_init:false`
|
||||
**Статус (2026-09-15):** репозиторий `kpa39l/email-assistant` создан на gitverse
|
||||
(private, id 336792, https://gitverse.ru/kpa39l/email-assistant); remote `gitverse`
|
||||
добавлен; master запушен (HEAD 39df85b51cd4263b0b3350e79eeb8926c7a95ec7).
|
||||
Пуш подтверждён ls-remote. Грядёт настройка pull mirror в gitea (см. project-git-setup).
|
||||
Пуш подтверждён ls-remote. Pull mirror в gitea `estorozhenko/email-assistant` создан
|
||||
(201, id 54, интервал 12h), HEAD совпадает с gitverse, mirror-sync 200.
|
||||
@@ -0,0 +1,7 @@
|
||||
# AGENTS.md
|
||||
|
||||
## Agent skills
|
||||
|
||||
- **Issue tracker**: Gitea — `hermes/email-assistant` on https://gitea.nixg.ru.
|
||||
See `docs/agents/issue-tracker.md` for the full workflow (tea CLI `gitea.nixg-full` login / gitea-mcp MCP tools, triage labels, blocking-edge convention).
|
||||
- **Domain docs**: single-context — `CONTEXT.md` and `docs/adr/` at repo root if/when they exist.
|
||||
@@ -1,7 +1,7 @@
|
||||
# Email Assistant — локальный архив и ассистент почты
|
||||
|
||||
**Дата:** 2026-09-13
|
||||
**Фаза:** 1.5–1.7 + Портфель веб-UI (планирование) + Классификация/обработчики (в работе)
|
||||
**Дата:** 2026-09-15
|
||||
**Фаза:** 1.5–1.7 + Портфель веб-UI (планирование) + Классификация/обработчики (в работе) + IMAP Realtime Sync (планирование)
|
||||
|
||||
**Стек:** Himalaya CLI → Python → SQLite → Ollama (Qwen3:8b) → Radicale (CalDAV) → Telegram
|
||||
|
||||
@@ -140,8 +140,46 @@ Hermes cron:
|
||||
- [x] **Фикс секретов (2026-09-14)**: скрипт теперь сам читает `radicale/.env` (RADICALE_PASS) и `/opt/vesti/.env` (VESTI_BOT_TOKEN) — раньше без ручного export был 401; добавлен stdlib-парсер .env (python-dotenv в системе нет)
|
||||
- [x] **Живой прогон (2026-09-14)**: письмо 2026/422 → VTODO «Задачи» (204) + VEVENT «Рабочий» (204); повтор — идемпотентно (0 дублей)
|
||||
- [x] **Cron (2026-09-14)**: `mail-classify-handlers` (6e1e78ceedfd, every 5m) — классификатор (--limit 10) → обработчики; end-to-end проверено: 5 новых «meeting» → 5 VEVENT (201)
|
||||
- [ ] Живое urgent-письмо → доставка в Telegram (механика готова, токен подхватывается; пока не было urgent-писем)
|
||||
- [x] **Telegram-секреты**: токен `VESTI_BOT_TOKEN` читается из /opt/vesti/.env; канал-дефолт `@dedinit_vesti` (TELEGRAM_CHAT_ID можно переопределить в .env проекта)
|
||||
- [x] **Живое urgent-письмо → доставка в Telegram (2026-09-15)**: письмо с id 3216 («RE: Платежи Аврора», Archive/2026/08/3216) помечено `classification: urgent` → `email_handlers.py` с `TELEGRAM_CHAT_ID=281328953` отправил в ЛС (private чат kpa39l) через бота @dedinit_controller_bot — `msg_id=60`, пользователь подтвердил получение; `handled_urgent: true` записан (идемпотентно, повторной отправки нет)
|
||||
- [x] **Telegram-секреты**: токен `VESTI_BOT_TOKEN` читается из /opt/vesti/.env; канал-дефолт `@dedinit_vesti` — **внимание**: для ЛС нужно явно `TELEGRAM_CHAT_ID=281328953` (private chat kpa39l); зафиксировать дефолт в ЛС — TODO при настройке stream-сервиса
|
||||
|
||||
---
|
||||
|
||||
## Сделано в сессии 2026-09-15 (IMAP realtime sync — планирование + Telegram-уведомления в ЛС)
|
||||
|
||||
- [x] **Выделена авторизация IMAP в отдельную функцию** — `scripts/imap_client.py` (новый): `imap_connect()` (socket+STARTTLS+LOGIN, re-try 5→60с, креды из config/himalaya-config.toml); `mail_archive.py` переведён на неё (sys.path-фикс для cron); проверено живьём: `OK: connected+LOGIN as e.storozhenko @ mail.corpoffice.tech`
|
||||
- [x] **Логирование авторизации + метрики** (в imap_client.py): 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_ok/err, sessions_active); контекст-менеджер `imap_session`; `--metrics` CLI
|
||||
- [x] **Telegram-уведомления о важных письмах — в ЛС** (kpa39l): создан `/opt/hermes/email-assistant/.env` с `TELEGRAM_CHAT_ID=281328953`; живой тест: письмо 3216 («RE: Платежи Аврора») помечено urgent → доставка в private chat 281328953 (bot @dedinit_controller_bot, msg_id=60), пользователь подтвердил
|
||||
- [x] **Найден и исправлен баг доставки**: крон-скрипт `mail-classify-handlers.sh` использовал системный `python3` без httpx → все Telegram-уведомления молча падали (накопилось 10 urgent-писем). Фикс: `PY=/opt/vesti/.venv/bin/python` + `--limit 2` (дозированная отправка backlog, анти-спам)
|
||||
- [x] **Формат даты в уведомлениях**: добавлены `format_date_for_tg()` (ISO → `02.09.2026 10:49`) и `_unquote_yaml()` (снятие YAML-кавычек в parse_email_md); строка `📅 ДД.ММ.ГГГГ ЧЧ:ММ` в handle_urgent
|
||||
- [x] **IMAP realtime sync (проект)**: proposal.md + design.md + tasks.md + spec-дельты (imap-realtime-sync, email-attachments, email-classification) — openspec `change 'imap-realtime-sync' is valid`; диагноз rate-limit (не TLS-fingerprinting); 2 микросервиса (imap_stream.py IDLE + change_analyzer.py ChangeLog SQLite), systemd, soft-delete
|
||||
- [x] **Backlog urgent-писем**: 10 писем в INBOX с `classification: urgent` без `handled_urgent` (июль 2026 — сентябрь 2026, включая срочное «переоформление договоров… отключение интернета до 20-го») — отправляются по 2 за крон-прогон в ЛС
|
||||
- [x] **Репозиторий переведён на gitverse (источник истины)**: создан приватный `kpa39l/email-assistant` на gitverse.ru (id 336792) через API; remote `gitverse` добавлен; master запушен; **gitea: pull mirror `estorozhenko/email-assistant` ← gitverse (12h, id 54)** — gitea.nixg.ru остаётся зеркалом, не источником
|
||||
- [x] **AGENT.md создан** — политика репозитория: gitverse = истина, gitea = зеркало, команды пуша, **местоположение токенов** (git-tokens.env, VESTI_BOT_TOKEN, TELEGRAM_CHAT_ID)
|
||||
|
||||
## Сделано в сессии 2026-09-15 (вечер) — imap_stream.py реализован + backlog urgent закрыт
|
||||
|
||||
### imap_stream.py — Фаза 1 (imap-realtime-sync, задачи 3–5) ✅
|
||||
- [x] **`scripts/imap_stream.py` (новый, ~690 строк)** — постоянный IMAP IDLE-поток на сыром socket поверх `imap_client.imap_connect` (aioimaplib не ставится):
|
||||
- IDLE-цикл: `idle_start` → `idle_wait(25с)` → `idle_done` → перевыпуск (Exchange рвёт IDLE ~60с, поэтому IDLE_TIMEOUT=25с)
|
||||
- Reconnect с паузой `CONNECT_PAUSE=45с` (rate-limit Exchange), цикл переживает обрыв (проверено live: 150с демон → обрыв → reconnect → SELECT)
|
||||
- `reconcile_new` — новые UID (архивация через mail_archive) + `reconcile_full` каждые 5 циклов (флаги/удаления)
|
||||
- SQLite `/opt/hermes/email/state/mailbox.db`: `mailbox_state` (uid, folder, message_id, in_reply_to, refs, flags, has_attachment, archive_path, last_seen, deleted) + `mailbox_events`
|
||||
- CLI: `--check`, `--test-idle`, `--status`, `--metrics`, без аргументов — демон
|
||||
- [x] **Live-тесты Exchange (пауза ≥45с между коннектами) — все прошли**: `--check` (LOGIN 365 писем), `--test-idle` (SELECT → reconcile → IDLE 30с → 0 событий → DONE, EXIT=0), демон 150с (обрыв+reconnect)
|
||||
- [x] **Исправлено live-тестами (питфолы, детали в WALKTHROUGH):** HERMES_REAL_HOME в imap_client (_load_credentials), `_cmd_ok` (хвост \r\n), колонка `references`→`refs` (SQLite), **UID-батчинг `1:1000` пуст** → `UID SEARCH ALL`/`UID <min>:*` (иначе все письма ложно deleted), Exchange не шлёт `+ idling` → `+ IDLE accepted`
|
||||
- [x] **Запушено в gitverse**: `5a81c86` (imap_stream + фикс токена TG + фильтр актуальности + HERMES_REAL_HOME)
|
||||
|
||||
### Backlog urgent (10 писем) — закрыт ✅
|
||||
- [x] **Корень найден**: `email_handlers.py` при наличии python-dotenv шёл в ветку `if load_dotenv` и НЕ грузил `/opt/vesti/.env` → `VESTI_BOT_TOKEN` пуст → все уведомления падали. Фикс: vesti/.env грузится всегда
|
||||
- [x] **5 сентябрьских писем доставлены в ЛС** (msg_id 62–64 и далее), помечены handled_urgent
|
||||
- [x] **Фильтр актуальности (запрос пользователя)**: `URGENT_MAX_AGE_DAYS=3` + `is_urgent_recent()` — письма старше 3 дней НЕ шлются в Telegram (помечаются handled_urgent); 5 старых (июль/авг) пропущены корректно
|
||||
- [x] Идемпотентность сохранена: дубли предотвращены (14229 помечен вручную)
|
||||
|
||||
### Открыто на следующую сессию
|
||||
- [ ] Запуск imap_stream.py как постоянного сервиса (systemd/cron) — код готов, не запущен навсегда
|
||||
- [ ] Баг Radicale: `[task] PUT 400: Bad Request` (VTODO-задача, не относится к urgent) — не разобран
|
||||
- [ ] Полный reconcile по всем подпапкам INBOX (сейчас Фаза 1 — только INBOX)
|
||||
|
||||
### Задача 1: Веб-интерфейс ассистента ⬜
|
||||
- [ ] FastAPI + SQLite FTS5: список писем (дата/адресант/тэги/папка)
|
||||
@@ -300,10 +338,10 @@ dav.example.com {
|
||||
|
||||
## Конфигурация
|
||||
|
||||
- **Репозиторий:** `https://gitea.nixg.ru/hermes/email-assistant`
|
||||
- **Push mirror на gitverse.ru:** ❌ не настроен
|
||||
- [ ] Настроить push mirror из gitea.nixg.ru в gitverse.ru
|
||||
- Требуется: создать репозиторий на gitverse.ru, получить токен, настроить mirror в настройках gitea (Settings → Git Hooks/Mirrors → Add Push Mirror)
|
||||
- **Репозиторий (источник истины):** `https://gitverse.ru/kpa39l/email-assistant` (private, remote `gitverse`, https с токеном из `/opt/hermes/.hermes/secrets/git-tokens.env`)
|
||||
- **Gitea (зеркало):** `https://gitea.nixg.ru/estorozhenko/email-assistant` (pull mirror ← gitverse, интервал 12h, id 54). Старый `hermes/email-assistant` — обычный репо, не зеркало.
|
||||
- **AGENT.md** — политика репозитория/токенов (см. файл)
|
||||
- ~~Push mirror gitea → gitverse~~ ✅ **сделано 2026-09-15** (pull mirror наоборот: gitverse → gitea)
|
||||
|
||||
### Ресурсы проекта (расположение и доступ)
|
||||
|
||||
|
||||
@@ -16,7 +16,7 @@
|
||||
| 2026-07-19 | clean_body — вырезание цитируемой переписки (Outlook/forward) | ✅ закрыта | STATUS.md контекст |
|
||||
| 2026-07-19 | Cron mail-index-incremental | 🔵 открыта | |
|
||||
| 2026-07-19 | Cron digest-weekly | 🔵 открыта | |
|
||||
| 2026-07-19 | Push mirror gitea → gitverse.ru | 🔵 открыта | |
|
||||
| 2026-07-19 | Push mirror gitea → gitverse.ru | ✅ закрыта (2026-09-15, иначе: pull mirror gitverse → gitea) | AGENT.md |
|
||||
|
||||
## 2026-09-11
|
||||
| Дата | Задача | Статус | Закрыта в |
|
||||
@@ -50,6 +50,23 @@
|
||||
| 2026-09-13 | Задача 2: Caddy reverse proxy — cal.nixg.ru работает (207), tasks.nixg.ru закомментирован | 🟡 частично (cal.nixg.ru готов) | STATUS.md |
|
||||
| 2026-09-13 | Задача 3: Android — контакты синхронизированы (DAVx5), события/задачи ещё не проверены в приложении | 🔵 открыта | |
|
||||
|
||||
## 2026-09-15
|
||||
| Дата | Задача | Статус | Закрыта в |
|
||||
|---|---|---|---|
|
||||
| 2026-09-15 | IMAP realtime sync: proposal/design/tasks/spec-дельты (openspec valid) | ✅ закрыта | openspec/changes/imap-realtime-sync/ |
|
||||
| 2026-09-15 | Выделена авторизация IMAP в `scripts/imap_client.py` (imap_connect, STARTTLS+LOGIN, re-try) + mail_archive.py переведён на неё | ✅ закрыта | STATUS.md §2026-09-15 |
|
||||
| 2026-09-15 | Логирование авторизации + метрики (JSON-лог /opt/hermes/email/logs/imap_client.log, imap_metrics, imap_session, --metrics) | ✅ закрыта | STATUS.md §2026-09-15 |
|
||||
| 2026-09-15 | Telegram-уведомления о важных письмах в ЛС (письмо 3216 → ЛС kpa39l, user подтвердил) | ✅ закрыта | `.env` проект (TELEGRAM_CHAT_ID=281328953) |
|
||||
| 2026-09-15 | Баг: крон mail-classify-handlers использовал python3 без httpx → уведомления не уходили (10 urgent-писем в backlog). Фикс: /opt/vesti/.venv/bin/python + --limit 2 | ✅ закрыта | scripts/mail-classify-handlers.sh |
|
||||
| 2026-09-15 | Формат даты в уведомлениях: `02.09.2026 10:49` (format_date_for_tg) + снятие YAML-кавычек (_unquote_yaml) | ✅ закрыта | scripts/email_handlers.py |
|
||||
| 2026-09-15 | imap_stream.py (IDLE-цикл) — РЕАЛИЗОВАН (Фаза 1, задачи 3–5): IDLE-цикл с reconnect (45с), reconcile new/full, SQLite mailbox.db (mailbox_state, mailbox_events); --check/--test-idle/--status/демон | ✅ закрыта | @session:default/20260915_174925_ec1ca4; git 5a81c86 |
|
||||
| 2026-09-15 | Backlog urgent: корень найден (vesti/.env не грузился при python-dotenv) → 5 сентябрьских доставлены; фильтр актуальности URGENT_MAX_AGE_DAYS=3 (старые не шлются) | ✅ закрыта | @session:default/20260915_174925_ec1ca4; git 5a81c86 |
|
||||
| 2026-09-15 | Фильтр неактуальных уведомлений (запрос пользователя): is_urgent_recent() — письмо старше 3 дн не шлётся, помечается handled_urgent | ✅ закрыта | scripts/email_handlers.py; git 5a81c86 |
|
||||
| 2026-09-15 | Live-тесты Exchange с паузой ≥45с — все прошли (--check, --test-idle 30с, демон 150с пережил обрыв+reconnect) | ✅ закрыта | @session:default/20260915_174925_ec1ca4 |
|
||||
| 2026-09-15 | Репозиторий: gitverse = источник истины (создан kpa39l/email-assistant, remote gitverse, master запушен) | ✅ закрыта | AGENT.md, STATUS.md |
|
||||
| 2026-09-15 | Gitea pull mirror: estorozhenko/email-assistant ← gitverse (12h, id 54) | ✅ закрыта | AGENT.md, STATUS.md |
|
||||
| 2026-09-15 | AGENT.md: добавить местоположение токенов (git-tokens.env, VESTI_BOT_TOKEN, TELEGRAM_CHAT_ID) | ✅ закрыта | AGENT.md |
|
||||
|
||||
## 2026-09-14
|
||||
| Дата | Задача | Статус | Закрыта в |
|
||||
|---|---|---|---|
|
||||
|
||||
+214
-1
@@ -2,6 +2,83 @@
|
||||
|
||||
Воспроизводимость: хронология, команды, решения, ошибки и как чинили.
|
||||
|
||||
## 2026-09-15 (вечер, сессия @session:default/20260915_174925_ec1ca4)
|
||||
|
||||
### imap_stream.py — Фаза 1 (IDLE-цикл) реализован
|
||||
|
||||
**Цель:** постоянный IMAP-поток (задачи 3–5 change imap-realtime-sync). aioimaplib
|
||||
не ставится → свой клиент на сыром socket поверх `imap_client.imap_connect`.
|
||||
|
||||
**Архитектура `scripts/imap_stream.py`:**
|
||||
- `ImapStream` — обёртка над сокетом: теги, `_cmd`, `_cmd_ok`, чтение до тега.
|
||||
- IDLE-цикл: `idle_start` → `idle_wait(25с)` → `idle_done` → перевыпуск.
|
||||
- Reconnect: при ошибке `CONNECT_PAUSE=45с` (rate-limit Exchange), повтор.
|
||||
- `reconcile_new(folder)` — поиск UID > last_uid, архивация через mail_archive,
|
||||
запись в `mailbox_state` + событие `mailbox_events`.
|
||||
- `reconcile_full(folder)` — полное сравнение UID/флагов (каждый 5-й цикл).
|
||||
- SQLite: `/opt/hermes/email/state/mailbox.db`, таблицы `mailbox_state` (uid,
|
||||
folder, message_id, in_reply_to, refs, flags, has_attachment, archive_path,
|
||||
last_seen, deleted) и `mailbox_events` (event/ts/folder/uid/detail).
|
||||
- CLI: `--check` (авторизация), `--test-idle` (IDLE 30с), `--status`, `--metrics`,
|
||||
без аргументов — демон.
|
||||
|
||||
**Питфолы Exchange, найденные live-тестами (пауза ≥45с между коннектами):**
|
||||
1. `_load_credentials` в imap_client.py читал `Path.home()` — в сессии Hermes
|
||||
HOME≠реальный. Фикс: `HERMES_REAL_HOME` (env) → `/home/estorozhenko`. Cron
|
||||
уже делает это в mail-archive.sh.
|
||||
2. `_cmd_ok`: `resp.split(b"\r\n")[-1]` давал пустой элемент (хвост `\r\n`) —
|
||||
команда считалась failed даже при OK. Фикс: разбор строк до тега.
|
||||
3. Колонка `references` — зарезервированное слово SQLite → `refs`.
|
||||
4. **UID-батчинг сломан:** `UID SEARCH UID 1:1000` пуст (UID — глобальный номер
|
||||
~14200+, не порядковый); первый пустой батч обрывал цикл, все письма ложно
|
||||
помечались deleted. Фикс: `UID SEARCH ALL` / `UID SEARCH UID <min>:*`
|
||||
(проверено: 365 UID).
|
||||
5. Exchange отвечает на IDLE `+ IDLE accepted, awaiting DONE command.` —
|
||||
НЕ `+ idling`. Фикс: матч `+ IDLE` или `+ idling`.
|
||||
6. **Exchange рвёт IDLE-соединение ~60с** (не 30 мин) — `IDLE_TIMEOUT=25с`
|
||||
(перевыпуск до серверного лимита). Демон переживает обрыв: reconnect 45с.
|
||||
|
||||
**Live-тесты (все с паузой ≥45с):**
|
||||
- `--check`: connect+LOGIN+SELECT 365 писем → OK.
|
||||
- `--test-idle`: SELECT → reconcile (0 новых, deleted=0, flags=18) → IDLE 30с →
|
||||
0 событий → DONE → EXIT=0.
|
||||
- Демон 150с: stream_started → SELECT 365 → обрыв IDLE (~60с) →
|
||||
reconnect_pause 45с → SELECT 365. Цикл переподключения работает.
|
||||
|
||||
**Git:** 5a81c86 запушен в gitverse (defcf53..5a81c86). Файлы: scripts/imap_stream.py
|
||||
(новый, ~690 строк), scripts/email_handlers.py (+фильтр актуальности), scripts/imap_client.py
|
||||
(+HERMES_REAL_HOME).
|
||||
|
||||
### Backlog urgent-pисем: корень найден + фильтр актуальности
|
||||
|
||||
**Проблема:** 10 urgent-писем не уходили в Telegram. Причина: при наличии
|
||||
python-dotenv `email_handlers.py` шёл в ветку `if load_dotenv` — грузил только
|
||||
`BASE_DIR/.env` и `radicale/.env`, НО НЕ `/opt/vesti/.env` → `VESTI_BOT_TOKEN`
|
||||
пуст → все обработки падали. Строка `✗ [urgent] ...: Нет токена Telegram`.
|
||||
|
||||
**Фикс (email_handlers.py):** `/opt/vesti/.env` грузится всегда (до ветвления).
|
||||
|
||||
**Запрос пользователя (mid-turn):** «перестань слать неактуальные срочные
|
||||
уведомления; встроить проверку на актуальность — сравнивать дату письма и
|
||||
текущую». Реализовано:
|
||||
- `URGENT_MAX_AGE_DAYS = 3` — письма старше 3 дней не шлются.
|
||||
- `is_urgent_recent(headers)` — парсит дату из frontmatter, сравнивает с now.
|
||||
- Старое письмо: `✓ [urgent] ...: пропущено (актуальность истекла)`, помечается
|
||||
`handled_urgent` (идемпотентность, не перебирается).
|
||||
|
||||
**Результат backlog:** 5 сентябрьских доставлены в ЛС (msg_id 62–64 и далее),
|
||||
5 старых (июль/авг) — пропущены фильтром корректно. Январьское «Сервер» — на
|
||||
деле `classification: info` (не urgent, ошибка подсчёта). 14229 помечен вручную
|
||||
(уже дважды уведомлён — предотвращён дубль).
|
||||
|
||||
**Известный отдельный баг (не urgent):** `✗ [task] ...: PUT 400: Bad Request`
|
||||
(Radicale VTODO-задача) — не разобран, открыт на следующую сессию.
|
||||
|
||||
### Cron-заметка
|
||||
|
||||
`mail-classify-handlers` (в venv /opt/vesti/.venv/bin/python, --limit 2/прогон)
|
||||
работает; после фикса токена обрабатывает по 2 письма за прогон идемпотентно.
|
||||
|
||||
## 2026-09-11
|
||||
|
||||
### Фаза 1.7: динамическое обнаружение подпапок INBOX
|
||||
@@ -303,4 +380,140 @@ VTIMEZONE Europe/Moscow + RRULE:FREQ=WEEKLY;BYDAY=TU, DTSTART 11:00.
|
||||
|
||||
**Cron:** mail-archive-every-5min (5f2305b2bbf8) приостанавливался на время
|
||||
отладки → **ВОЗОБНОВЛЁН** (next_run 08:05, state scheduled).
|
||||
вычитка openspec-файлов чейнджа.
|
||||
вычитка openspec-файлов чейнджа.
|
||||
|
||||
## 2026-09-15 — IMAP realtime sync (планирование) + Telegram-уведомления в ЛС
|
||||
|
||||
### 1. Диагностика авторизации: rate-limit, а не TLS-fingerprinting
|
||||
|
||||
**Симптом:** python/openssl/imaplib не логинятся на mail.corpoffice.tech:143
|
||||
(`NO AUTHENTICATE failed`), himalaya (rustls) логинится.
|
||||
|
||||
**Что перепробовано:**
|
||||
1. Raw-socket LOGIN для e.storozhenko / e.storozhenko@vinogorod.ru /
|
||||
VINOGOROD\e.storozhenko — NO.
|
||||
2. `ssl.wrap_socket` — `AttributeError` (убрано в py3.12) → `SSLContext.wrap_socket`.
|
||||
3. himalaya — логинится (эталон).
|
||||
4. Python с тем же base64 AUTHENTICATE PLAIN, что himalaya — NO.
|
||||
5. Сравнение TLS-отпечатков openssl vs rustls (ClientHello, ciphers) — различаются
|
||||
→ гипотеза **JA3 fingerprinting**.
|
||||
6. Порт 993 (IMAPS) через imaplib — тоже NO.
|
||||
7. **Опровержение:** чтение прод-кода `mail_archive.py` (~528) — он логинится
|
||||
простым `LOGIN login password` (не AUTH PLAIN). После паузы 45с — успех.
|
||||
Вывод: **rate-limit Exchange** (сервер молчит timeout после ~10 быстрых
|
||||
подключений), НЕ fingerprinting, НЕ бан IP.
|
||||
|
||||
**Урок:** не спешить с «экзотическими» диагнозами (JA3); сначала прочитать
|
||||
существующий прод-код — там уже рабочее решение.
|
||||
|
||||
### 2. `scripts/imap_client.py` — отдельная функция авторизации + наблюдаемость
|
||||
|
||||
- `imap_connect()`: socket → STARTTLS → LOGIN (простая форма, как в проде),
|
||||
re-try с экспоненциальной паузой 5→60с, креды из config/himalaya-config.toml.
|
||||
- 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_ok/err, sessions_active.
|
||||
- `imap_session` (contextmanager) — гарантирует session_ended.
|
||||
- CLI: `python3 imap_client.py` (self-test), `--metrics`.
|
||||
- `mail_archive.py`: блок соединения заменён на `imap_connect()`;
|
||||
`sys.path.insert(0, str(Path(__file__).parent))` перед импортом imap_client
|
||||
(важно для cron: иначе импорт падает из-за cwd).
|
||||
- Проверено: `HOME=/home/estorozhenko python3 imap_client.py` →
|
||||
`OK: connected+LOGIN as e.storozhenko @ mail.corpoffice.tech`.
|
||||
ВАЖНО: запускать с `HOME=/home/estorozhenko` (конфиг himalaya там).
|
||||
|
||||
### 3. Telegram-уведомления о важных письмах — в ЛС (kpa39l)
|
||||
|
||||
- Создан `/opt/hermes/email-assistant/.env`: `TELEGRAM_CHAT_ID=281328953`
|
||||
(private chat kpa39l, бот @dedinit_controller_bot id 7765665742).
|
||||
Файл в .gitignore (строка 9) — не попадёт в git.
|
||||
- Порядок загрузки .env в email_handlers.py: проект `.env` → radicale/.env →
|
||||
/opt/vesti/.env; `override=False` → проект не перезаписывается.
|
||||
- Живой тест: письмо 3216 («RE: Платежи Аврора») помечено
|
||||
`classification: urgent` → отправлено в ЛС (msg_id=60), пользователь
|
||||
подтвердил получение. `handled_urgent: true` записан.
|
||||
|
||||
### 4. БАГ: крон не отправлял уведомления (python3 без httpx)
|
||||
|
||||
**Симптом:** 10 urgent-писем накопились без уведомлений.
|
||||
|
||||
**Причина:** `scripts/mail-classify-handlers.sh` вызывал `python3` (системный),
|
||||
у которого нет httpx → `tg_call` бросал RuntimeError, скрипт глушил ошибку
|
||||
`|| echo "[handlers] ошибка"`.
|
||||
|
||||
**Фикс:** `PY=/opt/vesti/.venv/bin/python` (fallback python3) + `--limit 2`
|
||||
(дозированно, анти-спам backlog). Проверено: `bash -n` OK.
|
||||
|
||||
### 5. Формат даты в уведомлениях
|
||||
|
||||
- `format_date_for_tg()`: `2026-09-02 10:49+03:00` → `02.09.2026 10:49`
|
||||
(ДД.ММ.ГГГГ ЧЧ:ММ; RFC3339 и Z тоже; незнакомое — as-is).
|
||||
- `_unquote_yaml()` в parse_email_md: снятие YAML-кавычек (`"..."` → значение).
|
||||
- Строка `📅 <дата>` в handle_urgent под заголовком.
|
||||
|
||||
### 6. Репозиторий: gitverse = источник истины + gitea mirror
|
||||
|
||||
**Решение пользователя (2026-09-15):** gitverse.ru — источник истины для всех
|
||||
репозиториев; gitea.nixg.ru — зеркало.
|
||||
|
||||
**Токен:** единый источник — `/opt/hermes/.hermes/secrets/git-tokens.env`
|
||||
(chmod 600): `GITVERSE_LOGIN=kpa39l` (НЕ estorozhenko!), `GITVERSE_PAT`,
|
||||
`GITVERSE_API=https://api.gitverse.ru`, `GITEA_API=http://127.0.0.1:3000/api/v1`,
|
||||
`GITEA_TOKEN`. Загрузка: `set -a && source ... && set +a`.
|
||||
|
||||
**Создание репо на gitverse (private, auto_init:false!):**
|
||||
```bash
|
||||
curl -X POST "$GITVERSE_API/user/repos" \
|
||||
-H "Authorization: Bearer $GITVERSE_PAT" \
|
||||
-H "Accept: application/vnd.gitverse.object+json;version=1" \
|
||||
-H "Content-Type: application/json" \
|
||||
-d '{"name":"email-assistant","private":true,"auto_init":false}' # → 201
|
||||
```
|
||||
|
||||
**Push в gitverse:** remote `gitverse` = `https://kpa39l:$GITVERSE_PAT@gitverse.ru/kpa39l/email-assistant.git`; `git push gitverse master`.
|
||||
|
||||
**Pull mirror в gitea (12h):** migrate API:
|
||||
```bash
|
||||
curl -X POST "$GITEA_API/repos/migrate" -H "Authorization: token $GITEA_TOKEN" \
|
||||
-d '{"clone_addr":"https://oauth2:$GITVERSE_PAT@gitverse.ru/$GITVERSE_LOGIN/email-assistant.git",
|
||||
"repo_name":"email-assistant","repo_owner":"estorozhenko","service":"git",
|
||||
"mirror":true,"mirror_interval":"12h"}' # → 201
|
||||
```
|
||||
Проверка: `git ls-remote https://oauth2:$PAT@gitverse.ru/kpa39l/email-assistant.git HEAD` ==
|
||||
`docker exec gitea git --git-dir=/data/git/repositories/estorozhenko/email-assistant.git rev-parse HEAD`;
|
||||
mirror-sync `POST $GITEA_API/repos/estorozhenko/email-assistant/mirror-sync` → 200.
|
||||
|
||||
**Питфолы:**
|
||||
- База gitverse API — `api.gitverse.ru`, НЕ `gitverse.ru/api/v1` (там HTML SPA).
|
||||
- Gitea migrate только с `https://oauth2:<PAT>@...` (ssh:// → 422).
|
||||
- `auto_init:false` обязательно (иначе первый push упадёт).
|
||||
- `GITVERSE_API` уже содержит `https://` — не добавлять его повторно.
|
||||
- После пуша токен в remote URL — оставить (git config, не в коммитах); НЕ вычищать,
|
||||
иначе следующий push потребует креды.
|
||||
|
||||
### Открытые хвосты (2026-09-15)
|
||||
|
||||
- imap_stream.py (Фаза 1): скелет IDLE-цикла; aioimaplib не установлен —
|
||||
выбор: ставить его или свой клиент поверх imap_connect() (raw socket).
|
||||
- Backlog 10 urgent-писем (июль-сентябрь): отправляются по 2/прогон крона в ЛС;
|
||||
часть просрочена — решить, пропускать ли (пометить handled).
|
||||
|
||||
## 2026-09-16 — сессия закрыта (происшествие: путаница проектов)
|
||||
|
||||
Сессия ошибочно начата с выполнения задач **/opt/vesti/TODO.md**: агент неверно
|
||||
истолковал запрос «todo пустой?» и открыл чужой TODO (vesti), вместо своего
|
||||
(TODO.md email-assistant). Пользователь указал: **сессии email-assistant работают
|
||||
только с email-assistant**; задачи vesti живут в /opt/vesti/TODO.md и выполняются
|
||||
только по явной команде.
|
||||
|
||||
**Проверка (выполнена):** в файлах email-assistant (TODO.md, STATUS.md,
|
||||
WALKTHROUGH.md, AGENT.md) задач проекта vesti НЕТ — все упоминания vesti
|
||||
технические (VESTI_BOT_TOKEN из /opt/vesti/.env, python /opt/vesti/.venv/bin,
|
||||
дефолт канала @dedinit_vesti). Переносить нечего. Файлы email-assistant в этой
|
||||
сессии не изменялись.
|
||||
|
||||
**Побочный эффект:** в /opt/vesti остались незакоммиченные изменения от
|
||||
ошибочной работы (openspec/changes/manual-direction/, extract-post-sources/,
|
||||
sources/extract.py, правки web/app.py, web/templates/selected.html, TODO.md
|
||||
vesti, bot-флаг vesti на GtS=True) — судьбу решает пользователь.
|
||||
- Live-тесты Exchange: пауза ≥45с между подключениями (rate-limit).
|
||||
@@ -0,0 +1,69 @@
|
||||
# Issue Tracker
|
||||
|
||||
**Tracker:** Gitea (self-hosted, https://gitea.nixg.ru)
|
||||
**Repository:** `hermes/email-assistant`
|
||||
**Web UI:** https://gitea.nixg.ru/hermes/email-assistant
|
||||
|
||||
This is a custom ("Other") tracker configuration for the Matt Pocock engineering skills (`to-tickets`, `to-spec`, `triage`). Gitea is not natively supported by `/setup-matt-pocock-skills` (which offers GitHub / GitLab / local files), so the workflow is recorded here as freeform prose.
|
||||
|
||||
## How to publish issues
|
||||
|
||||
Use the **`tea` CLI** (v0.16.0, at `~/.local/bin/tea`) with the login **`gitea.nixg-full`**:
|
||||
|
||||
```bash
|
||||
export PATH="$HOME/.local/bin:$PATH"
|
||||
|
||||
# Create an issue
|
||||
tea issues create --login gitea.nixg-full --repo hermes/email-assistant \
|
||||
--title "Ticket title" \
|
||||
--description "Ticket body (use the to-tickets per-ticket template)" \
|
||||
--labels "ready-for-agent"
|
||||
|
||||
# List issues
|
||||
tea issues list --login gitea.nixg-full --repo hermes/email-assistant -o simple
|
||||
|
||||
# Show one issue (with comments)
|
||||
tea issues --login gitea.nixg-full --repo hermes/email-assistant 42 --comments
|
||||
|
||||
# Add a comment
|
||||
tea comment --login gitea.nixg-full --repo hermes/email-assistant 42 "comment"
|
||||
|
||||
# Close / reopen
|
||||
tea issues close --login gitea.nixg-full --repo hermes/email-assistant 42
|
||||
```
|
||||
|
||||
**Alternative:** the `gitea` MCP server (gitea-mcp, 72 tools: `create_issue`, `list_issues`, `update_issue`, `add_issue_labels`, ...) is configured in Hermes `mcp_servers.gitea` — it talks to the same instance with the same token. Use whichever is more convenient; prefer MCP tools when you are already in an agent session, `tea` for quick shell checks.
|
||||
|
||||
## Repository and credentials
|
||||
|
||||
- Instance: `https://gitea.nixg.ru` (API v1.26.2, external VPS 5.129.217.146)
|
||||
- Owner/repo: `hermes/email-assistant` (issues enabled: yes)
|
||||
- Login (tea): `gitea.nixg-full` — token `hermes-mcp-full` (full: `read:user`, `write:user`, `read:issue`, `write:issue`, `read:repository`, `write:repository`)
|
||||
- Token source: `/opt/hermes/.hermes/secrets/git-tokens.env` (`GITEA_NIXG_TOKEN`); also in Hermes `config.yaml` → `mcp_servers.gitea` env `GITEA_TOKEN`.
|
||||
- SSH (for git push): `ssh://git@gitea.nixg.ru:2222/hermes/email-assistant.git`
|
||||
|
||||
## Triage labels (created)
|
||||
|
||||
The five canonical labels are already created in this repo (IDs shown):
|
||||
|
||||
| Label | Color | ID |
|
||||
|---|---|---|
|
||||
| `needs-triage` | `#e11d21` | 1 |
|
||||
| `needs-info` | `#fbca04` | 2 |
|
||||
| `ready-for-agent` | `#0e8a16` | 3 |
|
||||
| `ready-for-human` | `#006b75` | 4 |
|
||||
| `wontfix` | `#ffffff` | 5 |
|
||||
|
||||
Use `ready-for-agent` for tickets published by `to-tickets` (agent-grabbable by construction, per the skill default).
|
||||
|
||||
## Blocking edges
|
||||
|
||||
Gitea has no native issue-blocking links in the API used here. Follow the `to-tickets` rule for trackers without native blocking: publish tickets in dependency order (blockers first), and set each ticket's **"Blocked by"** section to the blocking issues' real identifiers (e.g. `#12`, `#13`).
|
||||
|
||||
## Scope / local files fallback
|
||||
|
||||
If a task is not tied to `hermes/email-assistant` (e.g. infra on another repo), either:
|
||||
- create the issue in the matching repo on gitea.nixg.ru (`tea issues create --repo <owner>/<repo>`), or
|
||||
- fall back to local markdown under `.scratch/<feature>/issues/` in that project.
|
||||
|
||||
Update this file if the tracker repository ever changes.
|
||||
@@ -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-<folder>.json`.
|
||||
Уже есть: `fetch_attachments_imaplib()` (сырой IMAP, BODY.PEEK[], не ставит
|
||||
`\Seen`), `get_inbox_subfolders()` (динамическое обнаружение 136 папок,
|
||||
CHILDREN), `_imap_utf7_encode()` (modified UTF-7 для кириллических папок).
|
||||
- Архив: `/opt/hermes/email/<folder>/YYYY/MM/<uid>/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 (уже есть паттерн) |
|
||||
@@ -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`, классификатор, обработчики) не меняют поведение
|
||||
при отказе от сервисов — они работают как раньше.
|
||||
@@ -0,0 +1,27 @@
|
||||
# email-attachments Specification
|
||||
|
||||
## MODIFIED Requirements
|
||||
|
||||
### Requirement: Вложения сохраняются в каталог письма
|
||||
|
||||
Для каждого письма с вложениями (флаг `has_attachment: true` в frontmatter)
|
||||
вложения MUST быть сохранены в подкаталог `attachments/` каталога письма
|
||||
(`/opt/hermes/email/<folder>/YYYY/MM/<uid>/attachments/`). Скачивание выполняется
|
||||
как при опросе по расписанию (текущее поведение), так и при событийной
|
||||
синхронизации (сервис `imap-realtime-sync`). Скачивание вложений MUST NOT
|
||||
выставлять флаг `\Seen`.
|
||||
|
||||
#### Scenario: Письмо с вложением архивировано
|
||||
- **WHEN** `mail_archive.py` заархивировал письмо с `has_attachment: true`
|
||||
- **THEN** файлы вложений лежат в `<msg_dir>/attachments/` и совпадают с вложениями на IMAP-сервере
|
||||
|
||||
#### Scenario: Вложения при событийной синхронизации
|
||||
- **GIVEN** новое письмо с вложением обнаружено сервисом `email-imap-stream`
|
||||
- **WHEN** письмо архивируется по событию IDLE
|
||||
- **THEN** вложения скачаны в `<msg_dir>/attachments/`, флаг `\Seen` не выставлен
|
||||
|
||||
#### Scenario: Вложения при reconcile
|
||||
- **GIVEN** письмо с вложением было пропущено (обрыв соединения)
|
||||
- **WHEN** сервис выполняет reconcile
|
||||
- **THEN** вложения докачиваются в `<msg_dir>/attachments/` (то же поведение,
|
||||
что существующий `--attachments-backfill`)
|
||||
@@ -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** классификация не выполняется заново (идемпотентно, как сейчас)
|
||||
@@ -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` — архивация не
|
||||
останавливается полностью
|
||||
@@ -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 нет)
|
||||
@@ -74,11 +74,16 @@ if load_dotenv:
|
||||
load_dotenv(BASE_DIR / ".env", override=False)
|
||||
# radicale/.env — фактический источник RADICALE_PASS (проектного .env нет)
|
||||
load_dotenv(BASE_DIR / "radicale" / ".env", override=False)
|
||||
# Токен Telegram живёт в /opt/vesti/.env (проект-источник бота @dedinit_vesti).
|
||||
# Загружаем ВСЕГДА (и при наличии python-dotenv, и без него), иначе
|
||||
# VESTI_BOT_TOKEN никогда не попадает в env и все urgent-уведомления
|
||||
# падают с «Нет токена Telegram».
|
||||
vesti_env = Path("/opt/vesti/.env")
|
||||
if vesti_env.exists():
|
||||
load_dotenv(vesti_env, override=False)
|
||||
else:
|
||||
_load_env_file(BASE_DIR / ".env")
|
||||
_load_env_file(BASE_DIR / "radicale" / ".env")
|
||||
# Токен Telegram живёт в /opt/vesti/.env (проект-источник бота @dedinit_vesti);
|
||||
# опционально: если файл есть, берём VESTI_BOT_TOKEN/TELEGRAM_CHAT_ID оттуда.
|
||||
vesti_env = Path("/opt/vesti/.env")
|
||||
if vesti_env.exists():
|
||||
_load_env_file(vesti_env)
|
||||
@@ -103,11 +108,26 @@ TG_API = "https://api.telegram.org"
|
||||
LLM_TIMEOUT = 20
|
||||
MAX_PREVIEW_CHARS = 400 # превью письма для Telegram
|
||||
TG_TIMEOUT = 20
|
||||
# Письмо считается актуальным для urgent-уведомления, если оно не старше
|
||||
# этого срока (дни). Более старые НЕ шлём (пользователь просил не спамить
|
||||
# устаревшими срочными), но помечаем handled_urgent, чтобы не перебирать.
|
||||
URGENT_MAX_AGE_DAYS = 3
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 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 +140,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 +399,66 @@ 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 is_urgent_recent(headers):
|
||||
"""Актуально ли письмо для urgent-уведомления (не старше URGENT_MAX_AGE_DAYS).
|
||||
|
||||
Сравнивает дату письма (frontmatter date / internal_date) с текущей.
|
||||
Если дату не удалось распарсить — считаем актуальным (бить тревогу лучше,
|
||||
чем молчать). Возвращает bool.
|
||||
"""
|
||||
from datetime import datetime, timedelta
|
||||
s = (headers.get("date") or headers.get("internal_date") or "").strip()
|
||||
if not s:
|
||||
return True # нет даты — не фильтруем
|
||||
# форматы: '2026-09-02 10:49+03:00' | '2026-09-02T10:49:03+03:00'
|
||||
m = re.match(r"^(\d{4})-(\d{2})-(\d{2})", s)
|
||||
if not m:
|
||||
return True
|
||||
try:
|
||||
d = datetime(int(m.group(1)), int(m.group(2)), int(m.group(3)))
|
||||
except ValueError:
|
||||
return True
|
||||
limit = datetime.now() - timedelta(days=URGENT_MAX_AGE_DAYS)
|
||||
return d >= limit
|
||||
|
||||
|
||||
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 ""
|
||||
|
||||
# ФИЛЬТР АКТУАЛЬНОСТИ: старые письма не шлём (см. URGENT_MAX_AGE_DAYS).
|
||||
if not is_urgent_recent(headers):
|
||||
return True, f"пропущено: письмо старше {URGENT_MAX_AGE_DAYS} дн (актуальность истекла)"
|
||||
|
||||
preview = body.strip()
|
||||
if len(preview) > MAX_PREVIEW_CHARS:
|
||||
preview = preview[:MAX_PREVIEW_CHARS].rstrip() + "…"
|
||||
text = (
|
||||
f"⚠️ СРОЧНОЕ письмо\n\n"
|
||||
f"⚠️ СРОЧНОЕ письмо\n"
|
||||
f"{date_line}"
|
||||
f"От: {sender}\n"
|
||||
f"Тема: {subject}\n\n"
|
||||
f"{preview}\n\n"
|
||||
|
||||
@@ -0,0 +1,321 @@
|
||||
#!/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, чтобы иметь
|
||||
единый источник кред и не дублировать парсинг конфига.)
|
||||
|
||||
HOME: в сессии Hermes HOME может быть /opt/hermes/.hermes/home, а конфиг
|
||||
himalaya лежит в реальном HOME пользователя (/home/estorozhenko).
|
||||
Учитываем HERMES_REAL_HOME (как mail_archive._himalaya_cmd()).
|
||||
"""
|
||||
import tomllib
|
||||
from pathlib import Path
|
||||
|
||||
home = os.environ.get("HERMES_REAL_HOME") or str(Path.home())
|
||||
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)
|
||||
@@ -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())
|
||||
@@ -5,8 +5,14 @@ set -euo pipefail
|
||||
|
||||
cd /opt/hermes/email-assistant
|
||||
|
||||
# Классификатор: до 10 новых писем за запуск (Qwen ~10-20с/письмо = ~3 мин)
|
||||
python3 scripts/email_classifier.py --limit 10 || echo "[classifier] ошибка (продолжаем)" >&2
|
||||
# httpx (для Telegram в email_handlers.py) есть только в .venv /opt/vesti.
|
||||
# Системный python3 его не имеет — без этого уведомления не отправлялись.
|
||||
PY=/opt/vesti/.venv/bin/python
|
||||
[ -x "$PY" ] || PY=python3
|
||||
|
||||
# Обработчики: все письма с тегами, без handled_* (идемпотентно)
|
||||
python3 scripts/email_handlers.py || echo "[handlers] ошибка" >&2
|
||||
# Классификатор: до 10 новых писем за запуск (Qwen ~10-20с/письмо = ~3 мин)
|
||||
"$PY" scripts/email_classifier.py --limit 10 || echo "[classifier] ошибка (продолжаем)" >&2
|
||||
|
||||
# Обработчики: по 2 письма за проход (идемпотентно; избегаем залпа
|
||||
# уведомлений, если накопился backlog urgent)
|
||||
"$PY" scripts/email_handlers.py --limit 2 || echo "[handlers] ошибка" >&2
|
||||
+9
-22
@@ -480,26 +480,16 @@ def fetch_attachments_imaplib(uid, folder, dest_dir):
|
||||
import email as emailmod
|
||||
import email.header as email_header
|
||||
|
||||
try:
|
||||
creds = _himalaya_imap_credentials()
|
||||
except Exception as e:
|
||||
print(f" [WARN] нет IMAP-учётных из конфига himalaya: {e}", file=sys.stderr)
|
||||
return False
|
||||
|
||||
sock = None
|
||||
try:
|
||||
sock = socket.create_connection((creds["host"], creds["port"]), timeout=30)
|
||||
sock.settimeout(30)
|
||||
greet = sock.recv(1024)
|
||||
if not greet.startswith(b"* OK"):
|
||||
raise RuntimeError(f"bad greeting: {greet[:80]!r}")
|
||||
sock.sendall(b"a1 STARTTLS\r\n")
|
||||
resp = sock.recv(1024)
|
||||
if b"OK" not in resp:
|
||||
raise RuntimeError(f"STARTTLS failed: {resp[:80]!r}")
|
||||
ctx = sslmod.create_default_context()
|
||||
sock = ctx.wrap_socket(sock, server_hostname=creds["host"])
|
||||
sock.settimeout(30)
|
||||
# Авторизация — через общий imap_client.imap_connect() (отдельная функция,
|
||||
# с re-try против rate-limit Exchange). Креды читаются из конфига himalaya
|
||||
# внутри imap_connect(). Возвращает (sock, creds).
|
||||
import sys as _sys
|
||||
from pathlib import Path as _Path
|
||||
_sys.path.insert(0, str(_Path(__file__).resolve().parent))
|
||||
from imap_client import imap_connect
|
||||
sock, creds = imap_connect()
|
||||
|
||||
def cmd(tag, line, expect_literal=None):
|
||||
sock.sendall(f"{tag} {line}\r\n".encode())
|
||||
@@ -525,9 +515,6 @@ def fetch_attachments_imaplib(uid, folder, dest_dir):
|
||||
break
|
||||
return buf
|
||||
|
||||
r = cmd("a2", f'LOGIN {creds["login"]} {creds["password"]}')
|
||||
if b"OK LOGIN" not in r:
|
||||
raise RuntimeError(f"LOGIN failed: {r[-120:]!r}")
|
||||
# Папка — в IMAP modified UTF-7 (кириллица иначе не находится на Exchange)
|
||||
mbox = _imap_utf7_encode(folder)
|
||||
r = cmd("a3", f'SELECT "{mbox}"')
|
||||
@@ -589,7 +576,7 @@ def fetch_attachments_imaplib(uid, folder, dest_dir):
|
||||
pass
|
||||
|
||||
|
||||
def get_attachments(uid, folder, dest_dir, timeout=25):
|
||||
def get_attachments(uid, folder, dest_dir, timeout=120):
|
||||
"""
|
||||
Скачать вложения письма в dest_dir.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user