228 lines
8.5 KiB
Python
Executable file
228 lines
8.5 KiB
Python
Executable file
#!/usr/bin/env python3
|
||
import subprocess
|
||
import os
|
||
import sys
|
||
import logging
|
||
import requests
|
||
import json
|
||
import re
|
||
import psycopg2
|
||
from pathlib import Path
|
||
from datetime import datetime
|
||
|
||
# === CONFIG ===
|
||
BASE_DIR = Path("/home/matrixhasyou/mutt")
|
||
PYTHON_BIN = "/usr/bin/python3"
|
||
GOOGLE_PYTHON_BIN = "/home/matrixhasyou/swarm-services/domovoy-bot/venv-google/bin/python"
|
||
EXPORT_SCRIPT = BASE_DIR / "export_emails.py"
|
||
CONSOLIDATE_SCRIPT = BASE_DIR / "consolidate_to_mega_pdf.py"
|
||
DRIVE_SYNC_SCRIPT = Path("/home/matrixhasyou/domovoy_drive_sync.py")
|
||
MASTER_PDF = BASE_DIR / "ALL_EMAILS_CONSOLIDATED.pdf"
|
||
|
||
BOT_TOKEN = "8725618164:AAH1tGalq-pw1l0t4P5c0sdLkCVJdh7IE7M"
|
||
CHAT_ID = "197957361"
|
||
|
||
STATUS_FILE = BASE_DIR / "pipeline_status.json"
|
||
|
||
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
|
||
log = logging.getLogger("ChameleonSyncMaster")
|
||
|
||
def send_tg(text):
|
||
url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendMessage"
|
||
try:
|
||
requests.post(url, json={"chat_id": CHAT_ID, "text": text, "parse_mode": "HTML"},
|
||
timeout=20)
|
||
except Exception as e:
|
||
log.error(f"TG notification failed: {e}")
|
||
|
||
def get_db_url():
|
||
env_path = Path("/home/matrixhasyou/swarm-services/domovoy-bot/.env")
|
||
if env_path.exists():
|
||
try:
|
||
with open(env_path) as f:
|
||
for line in f:
|
||
if line.startswith("DATABASE_URL_SYNC="):
|
||
return line.split("=", 1)[1].strip()
|
||
except:
|
||
pass
|
||
return "postgresql://gemini_admin:secure_swarm_pass_2026@192.168.10.105:5433/domovoy_db"
|
||
|
||
def update_status(status, step, new_emails=0, total_emails=0, pdf_size_mb=0.0):
|
||
try:
|
||
data = {
|
||
"status": status,
|
||
"step": step,
|
||
"updated_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
|
||
"new_emails": new_emails,
|
||
"total_emails": total_emails,
|
||
"pdf_size_mb": round(pdf_size_mb, 2)
|
||
}
|
||
with open(STATUS_FILE, "w") as f:
|
||
json.dump(data, f)
|
||
except Exception as e:
|
||
log.error(f"Failed to write pipeline status: {e}")
|
||
|
||
def update_db_start():
|
||
try:
|
||
conn = psycopg2.connect(get_db_url())
|
||
cur = conn.cursor()
|
||
cur.execute("INSERT INTO email_export_logs (status, start_time, new_emails_count) VALUES (%s, CURRENT_TIMESTAMP, 0) RETURNING id", ('running',))
|
||
log_id = cur.fetchone()[0]
|
||
conn.commit()
|
||
cur.close()
|
||
conn.close()
|
||
return log_id
|
||
except Exception as e:
|
||
log.error(f"Failed to log start in DB: {e}")
|
||
return None
|
||
|
||
def update_db_finish(log_id, status, new_emails_count):
|
||
if not log_id: return
|
||
try:
|
||
conn = psycopg2.connect(get_db_url())
|
||
cur = conn.cursor()
|
||
cur.execute("""
|
||
UPDATE email_export_logs
|
||
SET status = %s, end_time = CURRENT_TIMESTAMP, new_emails_count = %s
|
||
WHERE id = %s
|
||
""", (status, new_emails_count, log_id))
|
||
conn.commit()
|
||
cur.close()
|
||
conn.close()
|
||
except Exception as e:
|
||
log.error(f"Failed to update log status in DB: {e}")
|
||
|
||
def update_email_audit_stats():
|
||
try:
|
||
all_folder = Path("/home/matrixhasyou/mutt/exported_emails/all")
|
||
if not all_folder.exists():
|
||
return
|
||
|
||
stats = {}
|
||
for f in all_folder.glob("*.pdf"):
|
||
match = re.match(r'([a-zA-Z0-9.-]+)_(\d{4}-\d{2}-\d{2})_(.*)', f.stem)
|
||
if match:
|
||
domain, dt_str, subject = match.group(1), match.group(2), match.group(3)
|
||
if domain not in stats:
|
||
stats[domain] = {"inbound": 0, "outbound": 0, "last_date": dt_str}
|
||
else:
|
||
if dt_str > stats[domain]["last_date"]:
|
||
stats[domain]["last_date"] = dt_str
|
||
|
||
# Простейшая эвристика направления
|
||
stats[domain]["inbound"] += 1
|
||
|
||
conn = psycopg2.connect(get_db_url())
|
||
cur = conn.cursor()
|
||
|
||
# Сначала очистим старое для актуализации
|
||
cur.execute("DELETE FROM email_audit")
|
||
|
||
for domain, s in stats.items():
|
||
cur.execute("""
|
||
INSERT INTO email_audit (domain, description, inbound_count, outbound_count, last_email_date, updated_at)
|
||
VALUES (%s, %s, %s, %s, %s, CURRENT_TIMESTAMP)
|
||
""", (domain, f"Ведомство {domain}", s["inbound"], s["outbound"], datetime.strptime(s["last_date"], "%Y-%m-%d")))
|
||
|
||
conn.commit()
|
||
cur.close()
|
||
conn.close()
|
||
log.info("Email audit stats updated successfully in PostgreSQL.")
|
||
except Exception as e:
|
||
log.error(f"Failed to update email_audit stats: {e}")
|
||
|
||
def run_step(name, cmd):
|
||
log.info(f"--- Step: {name} ---")
|
||
result = subprocess.run(cmd, capture_output=True, text=True)
|
||
if result.returncode != 0:
|
||
log.error(f"Step {name} FAILED!")
|
||
log.error(result.stderr)
|
||
return False, result.stderr
|
||
log.info(f"Step {name} completed.")
|
||
return True, result.stdout
|
||
|
||
def main():
|
||
start_time = datetime.now()
|
||
log.info(f"Starting Master Sync at {start_time}")
|
||
|
||
log_id = update_db_start()
|
||
|
||
# 1. Export Emails
|
||
update_status("running", "Сбор почты (Gmail Export)")
|
||
ok, out_export = run_step("Export Emails", [PYTHON_BIN, str(EXPORT_SCRIPT)])
|
||
if not ok:
|
||
update_status("error", "Ошибка на шаге: Сбор почты")
|
||
update_db_finish(log_id, "error", 0)
|
||
send_tg(f"❌ <b>Mutt Master Sync: FAILED</b>\nStep: Export Emails\nError: <code>{out_export[:200]}</code>")
|
||
return
|
||
|
||
# 2. Consolidate PDF
|
||
update_status("running", "Консолидация и разбиение PDF")
|
||
ok, out_consolidate = run_step("Consolidate PDF", [PYTHON_BIN, str(CONSOLIDATE_SCRIPT)])
|
||
if not ok:
|
||
update_status("error", "Ошибка на шаге: Консолидация PDF")
|
||
update_db_finish(log_id, "error", 0)
|
||
send_tg(f"❌ <b>Mutt Master Sync: FAILED</b>\nStep: Consolidate PDF\nError: <code>{out_consolidate[:200]}</code>")
|
||
return
|
||
|
||
# 3. Sync to G-Drive
|
||
update_status("running", "Синхронизация с Google Drive")
|
||
ok, out_drive = run_step("G-Drive Sync", [GOOGLE_PYTHON_BIN, str(DRIVE_SYNC_SCRIPT), str(MASTER_PDF)])
|
||
if not ok:
|
||
update_status("error", "Ошибка на шаге: Синхронизация Google Drive")
|
||
update_db_finish(log_id, "error", 0)
|
||
send_tg(f"❌ <b>Mutt Master Sync: FAILED</b>\nStep: G-Drive Sync\nError: <code>{out_drive[:200]}</code>")
|
||
return
|
||
|
||
end_time = datetime.now()
|
||
|
||
# Parsing Stats
|
||
new_emails = 0
|
||
total_emails = 0
|
||
for line in out_export.split('\n'):
|
||
if "Processed this run:" in line:
|
||
try: new_emails = int(line.split(':')[-1].strip())
|
||
except: pass
|
||
|
||
for line in out_consolidate.split('\n'):
|
||
if "Found" in line and "PDFs" in line:
|
||
try: total_emails = int(line.split('Found')[-1].split('PDFs')[0].strip())
|
||
except: pass
|
||
|
||
# Get total size of all parts or master pdf
|
||
pdf_size_mb = 0.0
|
||
if MASTER_PDF.exists():
|
||
pdf_size_mb = os.path.getsize(MASTER_PDF) / (1024 * 1024)
|
||
else:
|
||
parts = list(BASE_DIR.glob("ALL_EMAILS_CONSOLIDATED_part*.pdf"))
|
||
if parts:
|
||
pdf_size_mb = sum(p.stat().st_size for p in parts) / (1024 * 1024)
|
||
|
||
# Обновляем БД
|
||
update_db_finish(log_id, "success", new_emails)
|
||
update_email_audit_stats()
|
||
|
||
update_status("success", "Синхронизация завершена успешно", new_emails, total_emails, pdf_size_mb)
|
||
|
||
summary = f"""
|
||
🚀 <b>AI ПОЧТОВЫЙ ТУРБО-СИНХРОН: OK</b> 🚀
|
||
|
||
📥 <b>Новых писем:</b> <code>{new_emails}</code>
|
||
📚 <b>Всего в архиве:</b> <code>{total_emails}</code>
|
||
🐘 <b>Размер PDF:</b> <code>{pdf_size_mb:.2f} МБ</code>
|
||
☁️ <b>G-Drive:</b> Обновлено ✅
|
||
|
||
📅 {end_time.strftime("%Y-%m-%d %H:%M")}
|
||
<i>Робот закончил работу. NotebookLM готов к анализу.</i>
|
||
"""
|
||
send_tg(summary)
|
||
|
||
log.info("Master Sync completed successfully.")
|
||
|
||
# Save last sync time for web UI
|
||
sync_time_file = Path("/home/matrixhasyou/mutt/last_ai_sync.txt")
|
||
with open(sync_time_file, "w") as f:
|
||
f.write(end_time.strftime("%Y-%m-%d %H:%M:%S"))
|
||
|
||
if __name__ == "__main__":
|
||
main()
|