Files
email-assistant/scripts/contacts_extractor.py
T

594 lines
25 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
"""
contacts_extractor.py — извлечение адресной книги из подписей писем.
Сканирует INBOX-письма (email.md), вызывает Qwen3:8b через Ollama API
для парсинга подписи, дедуплицирует и пишет:
/opt/hermes/email/contacts/contacts.json — машиночитаемая база
/opt/hermes/email/contacts/index.json — email → contact_id
/opt/hermes/email/contacts/contacts.vcf — vCard 4.0 для импорта
/opt/hermes/email/contacts/last_scan.json — трекинг обработанных
"""
import hashlib
import json
import os
import re
import sys
import time
import urllib.error
import urllib.request
from datetime import date, datetime
from pathlib import Path
# ─── Конфиг ───────────────────────────────────────────────────────────────────
EMAIL_ROOT = Path("/opt/hermes/email")
CONTACTS_DIR = EMAIL_ROOT / "contacts"
STATE_DIR = EMAIL_ROOT / "state"
DB_PATH = EMAIL_ROOT / "mail_index.db"
OLLAMA_URL = "http://localhost:11434/api/generate"
OLLAMA_MODEL = "qwen3:8b"
# Папки, которые сканируем (только входящие)
SCAN_FOLDERS = ["INBOX", "INBOX/!Scan", "INBOX/!Битрикс", "INBOX/!ВГ Чек листы",
"INBOX/!Документооборот", "INBOX/!Завки", "INBOX/!Материалы",
"INBOX/!Отчеты", "INBOX/!Персонал", "INBOX/!Протоколы",
"INBOX/!Реестр оплаты", "INBOX/!Торик", "INBOX/Бюджет",
"INBOX/Контрагенты", "INBOX/ЛНД", "INBOX/Организация работы",
"INBOX/Приемка и стройка", "INBOX/Системы", "INBOX/Эксплуатация"]
LLM_TIMEOUT = 30 # секунд на один запрос к Ollama
MAX_BODY_CHARS = 5000 # обрезаем body для LLM (первые N символов)
# ─── Промпт ───────────────────────────────────────────────────────────────────
PROMPT_TEMPLATE = """Ты — экстрактор контактных данных из писем. Твоя задача — найти подпись ОТПРАВИТЕЛЯ (автора этого письма) и извлечь структурированные данные.
ВАЖНО: Тебе передаётся ТОЛЬКО новое сообщение, без цитируемой переписки. Подпись отправителя — в самом конце этого текста.
Правила поиска подписи:
- Подпись обычно отделена от тела письма: "С уважением,", "С наилучшими пожеланиями,", "Best regards,", "Kind regards,", "С ув.,", "—\\n", "—\\n", "—\\n"
- Если разделителя нет — последние 5-10 строк это подпись
- Ignore boilerplate (disclaimers, confidentiality notices)
Извлеки из подписи:
1. full_name — полное имя (ФИО) — только ОТПРАВИТЕЛЯ, не перепутай
2. email — email адрес (если есть в подписи, иначе null)
3. phone — основной телефон (в международном или местном формате)
4. phone_secondary — дополнительный телефон (если есть)
5. position — должность
6. company — название компании/организации
7. address — почтовый/юридический адрес (если есть)
8. raw_signature — полный текст найденной подписи (для отладки)
Если никакой подписи не найдено — верни JSON со всеми полями null.
Верни ТОЛЬКО JSON, без пояснений:
{{"full_name": null, "email": null, "phone": null, "phone_secondary": null, "position": null, "company": null, "address": null, "raw_signature": null}}
Вот текст письма:
{body}"""
# ─── Вспомогательные ──────────────────────────────────────────────────────────
def contact_id(email_addr):
"""MD5-хэш email для id контакта."""
if not email_addr:
return None
return hashlib.md5(email_addr.strip().lower().encode()).hexdigest()[:12]
def load_json(path, default=None):
"""Загрузить JSON, вернуть default если нет или битый."""
if default is None:
default = {}
try:
with open(path, "r", encoding="utf-8") as f:
return json.load(f)
except (FileNotFoundError, json.JSONDecodeError):
return default
def save_json(path, data):
"""Сохранить JSON атомарно."""
tmp = path.with_suffix(".tmp")
with open(tmp, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
tmp.replace(path)
def find_email_md_files(root):
"""Рекурсивно найти все email.md в папке."""
return sorted(root.rglob("email.md"))
def parse_email_md(path):
"""Прочитать email.md, вернуть (headers_dict, body_text)."""
content = path.read_text(encoding="utf-8", errors="replace")
# YAML frontmatter: первая строка "---", потом YAML, потом "---"
match = re.match(r"^---\s*\n(.*?)\n---\s*\n(.*)", content, re.DOTALL)
if match:
yaml_block = match.group(1)
body = match.group(2).strip()
# Парсим YAML вручную (без PyYAML)
headers = parse_simple_yaml(yaml_block)
else:
headers = {}
body = content.strip()
return headers, body
def parse_simple_yaml(text):
"""Примитивный парсер YAML frontmatter (ключ-значение, без вложенности)."""
result = {}
for line in text.strip().split("\n"):
m = re.match(r"^(\w[\w_-]*)\s*:\s*(.*)", line)
if m:
key = m.group(1)
val = m.group(2).strip().strip('"').strip("'")
result[key] = val
return result
def clean_body(body):
"""Очистить тело письма от HTML, цитируемой переписки и мусора."""
# Удаляем <#part ...> блоки
body = re.sub(r'<#part[^>]*>', '', body)
body = re.sub(r'<#/part>', '', body)
# Удаляем HTML-теги
body = re.sub(r'<[^>]+>', '', body)
# Удаляем mailto: ссылки
body = re.sub(r'\(mailto:[^)]+\)', '', body)
# Заменяем unicode-пробелы на обычные
body = re.sub(r'[\u00a0\u2000-\u200f\u2028-\u202f\u2060]+', ' ', body)
# Удаляем трекинг-ссылки
body = re.sub(r'https?://tn-eoc\.[^\s]+', '', body)
body = re.sub(r'https?://[^\s]+\?utm_[^\s]+', '', body)
# Удаляем цитируемую переписку — отрезаем всё от САМОГО РАННЕГО маркера цитирования
# Маркеры: Outlook (рус/англ headers), forwarded, > lines, андерскор-разделители
quote_patterns = [
r'^[\s]*_{4,}\s*$', # _____
r'От:.*\n[\s]*Отправлено:', # Russian Outlook (в любом месте строки)
r'^[\s]*From:.*\n[\s]*Sent:', # English Outlook headers
r'——-.*Forwarded.*——-',
r'——-.*Пересылаемое.*——-',
r'——-.*Original Message.*——-',
r'>.*\bwrote:',
]
earliest = None
earliest_pos = len(body)
for qp in quote_patterns:
for m in re.finditer(qp, body, re.MULTILINE):
if m.start() < earliest_pos:
earliest_pos = m.start()
earliest = m
if earliest:
body = body[:earliest_pos].strip()
else:
# Fallback: ищем любой маркер цитирования в последних 500 символах
# (От: Стороженко, From: ... — вшитые в строку подписи маркеры)
tail = body[-500:] if len(body) > 500 else body
for pattern in [r'От:', r'Отправлено:', r'From:', r'Sent:', r'Кому:', r'To:', r'Тема:', r'Subject:']:
m2 = re.search(pattern, tail)
if m2:
offset = len(body) - len(tail) + m2.start()
body = body[:offset].strip()
break
# Удаляем строки начинающиеся с > (если остались)
lines = body.split('\n')
cleaned = [l for l in lines if not re.match(r'^\s*>', l)]
body = '\n'.join(cleaned)
body = re.sub(r'\n{3,}', '\n\n', body)
return body.strip()
def call_llm(body_text, max_retries=2):
"""Вызвать Qwen через Ollama API, вернуть JSON."""
# Очищаем body
body_text = clean_body(body_text)
# Обрезаем body
truncated = body_text[:MAX_BODY_CHARS]
prompt = PROMPT_TEMPLATE.format(body=truncated)
for attempt in range(max_retries + 1):
if attempt > 0:
time.sleep(1)
payload = json.dumps({
"model": OLLAMA_MODEL,
"prompt": prompt,
"stream": False,
"options": {
"temperature": 0.1,
"num_predict": 1024,
}
}).encode("utf-8")
req = urllib.request.Request(
OLLAMA_URL,
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
try:
resp = urllib.request.urlopen(req, timeout=LLM_TIMEOUT)
data = json.loads(resp.read().decode("utf-8"))
response_text = data.get("response", "").strip()
except (urllib.error.URLError, json.JSONDecodeError, TimeoutError) as e:
if attempt < max_retries:
continue
print(f" ⚠ LLM error: {e}", file=sys.stderr)
return None
if not response_text:
if attempt < max_retries:
continue
print(f" ⚠ LLM empty response", file=sys.stderr)
return None
# Пробуем распарсить весь ответ как JSON
try:
return json.loads(response_text)
except json.JSONDecodeError:
pass
# Если не получилось — ищем { ... } внутри
brace_depth = 0
json_start = None
for i, ch in enumerate(response_text):
if ch == '{':
if brace_depth == 0:
json_start = i
brace_depth += 1
elif ch == '}':
brace_depth -= 1
if brace_depth == 0 and json_start is not None:
try:
return json.loads(response_text[json_start:i+1])
except json.JSONDecodeError:
pass
json_start = None
if attempt < max_retries:
continue
print(f" ⚠ LLM JSON parse error, raw: {response_text[:300]}", file=sys.stderr)
return None
# ─── Основная логика ──────────────────────────────────────────────────────────
def get_uid_from_path(path):
"""Извлечь UID из пути INBOX/YYYY/MM/UID/email.md."""
parts = path.parts
try:
# Ищем часть, которая является числом (UID)
for p in parts:
if p.isdigit():
return int(p)
except (ValueError, IndexError):
pass
return None
def get_folder_from_path(path, root):
"""Извлечь имя папки (INBOX/!Протоколы) относительно корня почты."""
rel = path.relative_to(root)
parts = rel.parts
# Формат: <folder>/YYYY/MM/UID/email.md
# folder может быть "INBOX" или "INBOX/!Протоколы"
folder_parts = []
for p in parts:
if p.isdigit() or re.match(r"^\d{4}$", p):
break
folder_parts.append(p)
return "/".join(folder_parts)
def get_date_from_path(path):
"""Извлечь дату из пути (по году/месяцу)."""
parts = path.parts
year = None
month = None
for p in parts:
if re.match(r"^\d{4}$", p) and 2020 <= int(p) <= 2030:
year = int(p)
elif re.match(r"^\d{2}$", p) and 1 <= int(p) <= 12:
month = int(p)
if year and month:
return date(year, month, 1)
return date.today()
def save_progress(contacts_db, contacts_dir, processed_uids, processed_emails):
"""Инкрементальное сохранение контактов и last_scan."""
contacts_db["contacts"].sort(key=lambda c: c.get("email", ""))
save_json(contacts_dir / "contacts.json", contacts_db)
save_json(contacts_dir / "index.json", contacts_db.get("by_email", {}))
generate_vcard(contacts_dir, contacts_db["contacts"])
last_scan = {
"last_processed": datetime.now().isoformat(timespec="seconds"),
"processed_uids": processed_uids,
"processed_emails": sorted(processed_emails),
}
save_json(contacts_dir / "last_scan.json", last_scan)
print(f" 💾 Сохранено ({len(contacts_db['contacts'])} контактов)", flush=True)
def scan(limit=0):
"""Основной цикл сканирования с SQLite-трекингом."""
import sqlite3
contacts_dir = CONTACTS_DIR
contacts_dir.mkdir(parents=True, exist_ok=True)
# Загружаем контакты
contacts_db = load_json(contacts_dir / "contacts.json", {"version": 1, "contacts": [], "by_email": {}})
today_str = date.today().isoformat()
# Открываем SQLite
if not DB_PATH.exists():
print("❌ mail_index.db не найден. Сначала запусти: python3 mail_index.py")
return
conn = sqlite3.connect(str(DB_PATH))
conn.row_factory = sqlite3.Row
# Берём письма, где контакты ещё не извлечены
cursor = conn.cursor()
cursor.execute("""
SELECT rowid, path, uid, folder, from_addr, subject
FROM emails
WHERE contacts_extracted = 0 AND contacts_skipped = 0
ORDER BY folder, uid
LIMIT ?
""", (limit if limit > 0 else 999999,))
pending = cursor.fetchall()
if not pending:
print(" Нет новых писем для обработки")
conn.close()
return
print(f" Найдено писем к обработке: {len(pending)}", flush=True)
total_new = 0
total_skipped = 0
processed_count = 0
save_counter = 0
for row in pending:
rowid = row["rowid"]
rel_path = row["path"]
uid = row["uid"]
folder = row["folder"]
sender = row["from_addr"]
subject = row["subject"]
file_path = EMAIL_ROOT / rel_path
if not file_path.exists():
# Письмо удалено — отмечаем чтоб не дёргать
cursor.execute("UPDATE emails SET contacts_skipped=1 WHERE rowid=?", (rowid,))
continue
# Извлекаем email отправителя
from_email = None
email_match = re.search(r'<([^>]+@[^>]+)>', sender)
if email_match:
from_email = email_match.group(1).strip().lower()
elif "@" in sender:
from_email = sender.strip().lower()
# Пропускаем, если уже обработан (по email)
if from_email and from_email in contacts_db["by_email"]:
contact = contacts_db["contacts"][contacts_db["by_email"][from_email]]
last_seen = contact.get("last_seen", "2000-01-01")
days_since = (date.today() - date.fromisoformat(last_seen)).days
if days_since < 30:
cursor.execute("UPDATE emails SET contacts_extracted=1, last_scanned=? WHERE rowid=?", (datetime.now().isoformat(timespec="seconds"), rowid))
conn.commit()
total_skipped += 1
continue
# Читаем тело письма
headers, body = parse_email_md(file_path)
print(f" UID {uid} ({sender[:50]})...", end=" ", flush=True)
# Вызываем LLM
llm_result = call_llm(body)
if llm_result is None:
print("⚠ skip (LLM error)", flush=True)
cursor.execute("UPDATE emails SET contacts_skipped=1, last_scanned=? WHERE rowid=?", (datetime.now().isoformat(timespec="seconds"), rowid))
conn.commit()
total_skipped += 1
continue
# Если контакт не найден
full_name = (llm_result.get("full_name") or "").strip()
if not full_name or full_name == "null":
print("→ нет подписи", flush=True)
cursor.execute("UPDATE emails SET contacts_skipped=1, last_scanned=? WHERE rowid=?", (datetime.now().isoformat(timespec="seconds"), rowid))
conn.commit()
total_skipped += 1
continue
# Извлекаем email из LLM-результата
contact_email = llm_result.get("email") or from_email
if contact_email and contact_email.strip().lower() != "null":
contact_email = contact_email.strip().lower()
else:
contact_email = from_email
if not contact_email:
print("→ нет email, skip", flush=True)
cursor.execute("UPDATE emails SET contacts_skipped=1, last_scanned=? WHERE rowid=?", (datetime.now().isoformat(timespec="seconds"), rowid))
conn.commit()
total_skipped += 1
continue
cid = contact_id(contact_email)
source_path = str(rel_path)
# Создаём запись контакта
new_contact = {
"id": cid,
"email": contact_email,
"full_name": full_name,
"phone": llm_result.get("phone") or None,
"phone_secondary": llm_result.get("phone_secondary") or None,
"position": llm_result.get("position") or None,
"company": llm_result.get("company") or None,
"address": llm_result.get("address") or None,
"first_seen": today_str,
"last_seen": today_str,
"source_uids": [source_path],
"source_folders": [folder],
}
# Дедупликация
existing_idx = contacts_db["by_email"].get(contact_email)
if existing_idx is not None:
existing = contacts_db["contacts"][existing_idx]
for key in ["full_name", "phone", "phone_secondary", "position", "company", "address"]:
if new_contact.get(key):
existing[key] = new_contact[key]
existing["last_seen"] = today_str
if source_path not in existing["source_uids"]:
existing["source_uids"].append(source_path)
if folder not in existing["source_folders"]:
existing["source_folders"].append(folder)
print(f"✓ обновлён: {full_name} <{contact_email}>", flush=True)
else:
contacts_db["contacts"].append(new_contact)
contacts_db["by_email"][contact_email] = len(contacts_db["contacts"]) - 1
print(f"✓ новый: {full_name} <{contact_email}>", flush=True)
# Отмечаем в SQLite
cursor.execute("UPDATE emails SET contacts_extracted=1, last_scanned=? WHERE rowid=?", (datetime.now().isoformat(timespec="seconds"), rowid))
conn.commit()
total_new += 1
processed_count += 1
save_counter += 1
if limit > 0 and processed_count >= limit:
print(f" ⏸ лимит {limit} достигнут", flush=True)
break
# Сохраняемся каждые 5 писем
if save_counter >= 5:
save_counter = 0
contacts_db["contacts"].sort(key=lambda c: c.get("email", ""))
save_json(contacts_dir / "contacts.json", contacts_db)
save_json(contacts_dir / "index.json", contacts_db.get("by_email", {}))
generate_vcard(contacts_dir, contacts_db["contacts"])
print(f" 💾 Сохранено ({len(contacts_db['contacts'])} контактов, uid={uid})", flush=True)
# Финальное сохранение
contacts_db["contacts"].sort(key=lambda c: c.get("email", ""))
save_json(contacts_dir / "contacts.json", contacts_db)
save_json(contacts_dir / "index.json", contacts_db.get("by_email", {}))
generate_vcard(contacts_dir, contacts_db["contacts"])
# Статистика
cursor.execute("SELECT COUNT(*) as cnt FROM emails WHERE contacts_extracted=1")
extracted_total = cursor.fetchone()["cnt"]
cursor.execute("SELECT COUNT(*) as cnt FROM emails WHERE contacts_skipped=1")
skipped_total = cursor.fetchone()["cnt"]
cursor.execute("SELECT COUNT(*) as cnt FROM emails WHERE contacts_extracted=0 AND contacts_skipped=0")
remaining = cursor.fetchone()["cnt"]
conn.close()
print(f"\n{'=' * 50}")
print(f"Готово. Новых контактов: {total_new}, пропущено: {total_skipped}")
print(f"Всего в базе: {len(contacts_db['contacts'])} контактов")
print(f"Индекс: {extracted_total} извлечено, {skipped_total} пропущено, {remaining} осталось")
def generate_vcard(contacts_dir, contacts):
"""Сгенерировать contacts.vcf (vCard 4.0)."""
lines = []
for c in contacts:
if not c.get("email"):
continue
lines.append("BEGIN:VCARD")
lines.append("VERSION:4.0")
lines.append(f"FN:{c['full_name']}")
# N:Фамилия;Имя;Отчество;;
name_parts = c["full_name"].split(maxsplit=2)
if len(name_parts) >= 2:
n_line = f"N:{name_parts[-1]};{name_parts[0]};{' '.join(name_parts[1:-1])};;"
else:
n_line = f"N:{c['full_name']};;;;"
lines.append(n_line)
lines.append(f"EMAIL;TYPE=WORK:{c['email']}")
if c.get("phone"):
lines.append(f"TEL;TYPE=WORK:{c['phone']}")
if c.get("phone_secondary"):
lines.append(f"TEL;TYPE=CELL:{c['phone_secondary']}")
if c.get("position"):
lines.append(f"TITLE:{c['position']}")
if c.get("company"):
lines.append(f"ORG:{c['company']}")
if c.get("address"):
lines.append(f"ADR;TYPE=WORK:;;{c['address']};;;")
# NOTE с источником
sources = ", ".join(c.get("source_uids", []))
lines.append(f"NOTE:Извлечено из писем: {sources}")
lines.append("END:VCARD")
lines.append("")
vcf_path = contacts_dir / "contacts.vcf"
vcf_path.write_text("\n".join(lines), encoding="utf-8")
print(f" vCard: {vcf_path} ({len(contacts)} контактов)")
# ─── CLI ──────────────────────────────────────────────────────────────────────
def main():
import argparse
parser = argparse.ArgumentParser(description="Извлечение контактов из писем")
parser.add_argument("--dry-run", action="store_true", help="Не вызывать LLM, только показать что будет обработано")
parser.add_argument("--limit", type=int, default=0, help="Максимум писем для обработки (0 = все)")
args = parser.parse_args()
print("📇 Contacts Extractor")
print(f" База: {CONTACTS_DIR}")
print(f" Модель: {OLLAMA_MODEL}")
if args.dry_run:
print(" 🔍 DRY RUN — без LLM\n")
# В dry-run просто сканируем
if args.dry_run:
contacts_dir = CONTACTS_DIR
contacts_dir.mkdir(parents=True, exist_ok=True)
last_scan = load_json(contacts_dir / "last_scan.json")
processed_uids = last_scan.get("processed_uids", {})
total_to_process = 0
for folder in SCAN_FOLDERS:
folder_path = EMAIL_ROOT / folder
if not folder_path.is_dir():
continue
processed_for_folder = processed_uids.get(folder, 0)
email_files = find_email_md_files(folder_path)
new_files = [f for f in email_files if (get_uid_from_path(f) or 0) > processed_for_folder]
print(f" {folder}: {len(email_files)} total, {len(new_files)} new (after UID {processed_for_folder})")
total_to_process += len(new_files)
print(f"\nВсего новых писем к обработке: {total_to_process}")
return
scan(limit=args.limit)
if __name__ == "__main__":
main()