gmao/ent_watchdog.py
root d10888c005
Some checks are pending
CI - Tests et Syntax / lint-and-test (push) Waiting to run
Corriger les erreurs OpenRouter et exposer Adminer
2026-08-16 13:22:14 +00:00

478 lines
19 KiB
Python

#!/usr/bin/env python3
"""Watchdog pour l'interprétation automatique des messages ENT - Pattern DND."""
import sys
import os
import time
import json
import signal
from datetime import datetime, timedelta, timezone
# Définir le répertoire de travail
WORKDIR = os.environ.get('DOCKER_GMAO_WORKDIR', '/root/hermes-workspace/gmao-college')
os.chdir(WORKDIR)
sys.path.insert(0, WORKDIR)
# Activer le venv
venv_site_packages = os.path.join(WORKDIR, 'venv', 'lib', 'python3.13', 'site-packages')
if os.path.exists(venv_site_packages):
sys.path.insert(0, venv_site_packages)
# Charger les variables d'environnement depuis .env
from dotenv import load_dotenv
env_file = os.environ.get('DOCKER_GMAO_ENVFILE', os.path.join(WORKDIR, '.env'))
load_dotenv(env_file)
# Fichier de log et de dernière exécution
LOG_FILE = '/tmp/ent_watchdog.log'
LAST_RUN_FILE = '/tmp/ent_last_run.txt'
running = True
def log(message, level='info'):
"""Écrit un message dans le log et dans la base centralisee."""
timestamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
print(f'[{timestamp}] {message}')
with open(LOG_FILE, 'a') as f:
f.write(f'[{timestamp}] {message}\n')
try:
sys.path.insert(0, '/app')
from app_new.watchdog_logger import log_watchdog
log_watchdog('ent_watchdog', message, level=level)
except Exception:
pass
def get_last_run():
"""Retourne la date de dernière exécution."""
if os.path.exists(LAST_RUN_FILE):
try:
with open(LAST_RUN_FILE, 'r') as f:
return datetime.fromisoformat(f.read().strip())
except:
return None
return None
def save_last_run():
"""Sauvegarde la date d'exécution."""
with open(LAST_RUN_FILE, 'w') as f:
f.write(datetime.now().isoformat())
def should_run(gmao_context):
"""Vérifie si le watchdog doit s'exécuter."""
if not gmao_context or not gmao_context.auto_interpret_enabled:
return False
last_run = get_last_run()
if not last_run:
return True
interval_minutes = gmao_context.auto_interpret_interval or 10
next_run = last_run + timedelta(minutes=interval_minutes)
return datetime.now() >= next_run
def sync_ent_messages():
"""Synchronise les messages ENT77."""
from app_new import create_app, db
from app_new.ent.models import EntCredential, EntMessage
from app_new.lib_ext.ent_service import login_ent, fetch_new_messages
from app_new.lib_ext.ent_crypto import decrypt_value
import requests
app = create_app()
with app.app_context():
# Récupérer les identifiants actifs
credentials = EntCredential.query.filter_by(is_active=True).all()
if not credentials:
log("Aucun identifiant ENT actif")
return {'success': False, 'error': 'Aucun identifiant ENT actif'}
total_saved = 0
total_updated = 0
for cred in credentials:
try:
# Déchiffrer les identifiants
username = decrypt_value(cred.ent_username_encrypted, app)
password = decrypt_value(cred.ent_password_encrypted, app)
# Connexion ENT
sess = requests.Session()
headers = login_ent(sess, username, password)
if not headers:
log(f"Échec de connexion ENT pour credential {cred.id}")
continue
# Récupérer les messages
messages = fetch_new_messages(sess, headers, cred.last_msg_id)
for msg in messages:
msg_id = str(msg.get('id'))
existing = EntMessage.query.filter_by(id=msg_id).first()
if not existing:
# Récupérer le détail du message
from app_new.lib_ext.ent_service import DETAIL_URL, extract_sender
detail_resp = sess.get(DETAIL_URL + msg_id, headers=headers, timeout=10)
if detail_resp.status_code == 200:
detail = detail_resp.json()
import re
body_text = re.sub('<[^<]+?>', '', detail.get('body', ''))
# Extraire l'expéditeur avec la fonction dédiée
sender_display, is_forwarded, original_sender = extract_sender(detail)
ts = detail.get('date')
msg_date = datetime.fromtimestamp(ts / 1000) if ts else datetime.utcnow()
# Récupérer le sender_id depuis le JSON
from_ = detail.get('from', {})
sender_id = str(from_.get('id', '')) if from_ else None
ent_msg = EntMessage(
id=msg_id,
credential_id=cred.id,
subject=detail.get('subject', 'N/A'),
sender_name=sender_display,
sender_id=sender_id,
date=msg_date,
body=body_text,
is_forwarded=is_forwarded,
original_sender=original_sender,
is_unread=msg.get('unread', False),
has_attachment=msg.get('hasAttachment', False),
attachment_count=msg.get('attachmentCount', 0)
)
db.session.add(ent_msg)
total_saved += 1
else:
total_updated += 1
# Mettre à jour le dernier message traité
if messages:
cred.last_msg_id = str(messages[0].get('id'))
cred.last_check = datetime.now(timezone.utc).replace(tzinfo=None)
except Exception as e:
import traceback
log(f"Erreur sync ENT credential {cred.id}: {str(e)}", level='error')
log(f"Traceback: {traceback.format_exc()[-500:]}", level='error')
db.session.commit()
log(f"Sync ENT: {total_saved} nouveaux, {total_updated} mis à jour")
return {'success': True, 'saved': total_saved, 'updated': total_updated}
def get_openrouter_model():
"""Recupere le nom du modele OpenRouter depuis les parametres."""
try:
sys.path.insert(0, '/app')
from app_new.core.models.settings import AppSettings
return AppSettings.get('openrouter_model', 'poolside/laguna-m.1:free')
except Exception:
return 'poolside/laguna-m.1:free'
def analyze_ent_message_with_hermes(msg, gmao_context):
"""Analyse un message ENT avec Hermes et retourne les interprétations."""
from app_new.core.models.settings import AppSettings
import requests
# AppSettings.get déchiffre les valeurs sensibles stockées sous la forme
# ``enc:...``. Lire directement ``.value`` envoyait la valeur chiffrée à
# OpenRouter et provoquait une réponse 401.
api_key = AppSettings.get('openrouter_api_key')
if not api_key:
return {'error': 'Clé API OpenRouter non configurée'}
# Préparer le contenu du message
msg_content = f"""Sujet: {msg.subject or '(sans sujet)'}
De: {msg.sender_name or 'Inconnu'}
Date: {msg.date.strftime('%d/%m/%Y %H:%M') if msg.date else '-'}
Contenu:
{msg.body or '(pas de contenu)'}
"""
# Construire le prompt
prompt = """Tu es un assistant pour la GMAO (Gestion de Maintenance Assistée par Ordinaire) d'un établissement scolaire.
Analyse ce message et EXTRAIS les informations pour créer une ou plusieurs interventions/tâches.
IMPORTANT: Tu DOIS répondre UNIQUEMENT avec un JSON valide, sans texte avant ni après.
Types possibles:
- "curative": Panne, problème, réparation nécessaire (ex: fuite, équipement cassé)
- "preventive": Maintenance planifiée, contrôle périodique (ex: inspection, entretien)
- "administrative": Tâche administrative à effectuer (ex: déposer des documents, remplir un formulaire, rappel)
- "formation": Formation à suivre, session de formation (ex: formation sécurité, recyclage)
- "information": Information simple, pas d'action requise
Actions possibles:
- "create": Créer une nouvelle intervention/tâche
- "update": Mettre à jour une intervention existante
- "inform": Juste informer, pas d'action
Format de réponse OBLIGATOIRE:
{
"analyses": [
{
"type": "curative|preventive|administrative|formation|information",
"action": "create|update|inform",
"title": "Titre court et précis de l'intervention",
"description": "Description détaillée du problème ou de la tâche",
"urgency": "haute|normale|basse",
"location": "Lieu concerné (salle, bâtiment)",
"equipment": "Équipement concerné (nom exact si identifiable, sinon 'Divers')",
"suggested_date": "YYYY-MM-DD ou null",
"notes": "Notes additionnelles"
}
]
}
Si plusieurs tâches distinctes sont mentionnées, crée plusieurs entrées dans "analyses".
Si c'est juste de l'information sans action, utilise "action": "inform".
"""
# Ajouter le contexte GMAO
if gmao_context:
full_context = gmao_context.get_full_context()
if full_context:
prompt += f"\n\n=== CONTEXTE GMAO ===\n{full_context}\n"
prompt += f"\n\n=== MESSAGE À ANALYSER ===\n{msg_content}\n=== FIN DU MESSAGE ===\n\nFournis ton analyse:"
# Appeler OpenRouter
try:
response = requests.post(
'https://openrouter.ai/api/v1/chat/completions',
headers={
'Authorization': f'Bearer {api_key}',
'Content-Type': 'application/json'
},
json={
'model': get_openrouter_model(),
'messages': [{'role': 'user', 'content': prompt}],
'max_tokens': 2000
},
timeout=60
)
if response.status_code != 200:
log(f"API erreur {response.status_code}: {response.text[:200]}")
return {'error': f'Erreur API: {response.status_code}', 'details': response.text[:200]}
data = response.json()
content = data.get('choices', [{}])[0].get('message', {}).get('content', '')
if not content:
log(f"API réponse vide: {data}")
return {'error': 'Réponse API vide'}
# Parser le JSON
try:
# Vérifier si le contenu est vide
if not content or not content.strip():
log(f"API contenu vide")
return {'type': 'information', 'action': 'inform', 'description': 'Pas de réponse API'}
log(f"API réponse: {content[:100]}...")
# Nettoyer le contenu (enlever les balises markdown)
if '```json' in content:
content = content.split('```json')[1].split('```')[0]
elif '```' in content:
content = content.split('```')[1].split('```')[0]
cleaned = content.strip()
try:
result = json.loads(cleaned)
except json.JSONDecodeError:
start = cleaned.find('{')
if start < 0:
raise
result, _ = json.JSONDecoder().raw_decode(cleaned[start:])
# Retourner la première analyse
if result.get('analyses') and len(result['analyses']) > 0:
analysis = result['analyses'][0]
return {
'type': analysis.get('type', 'information'),
'action': analysis.get('action', 'inform'),
'title': analysis.get('title', msg.subject),
'description': analysis.get('description', ''),
'urgency': analysis.get('urgency', 'normale'),
'location': analysis.get('location'),
'equipment': analysis.get('equipment'),
'confidence': 0.8
}
return {'type': 'information', 'action': 'inform', 'description': 'Pas d\'action requise'}
except json.JSONDecodeError as e:
return {'error': f'Erreur parsing JSON: {str(e)}', 'content': content[:200]}
except Exception as e:
return {'error': f'Erreur: {str(e)}'}
def run_watchdog():
"""Exécute un cycle du watchdog - retourne False si désactivé."""
global running
from app_new import create_app, db
from app_new.outlook.models import GmaoContext
from app_new.ent.models import EntMessage
from app_new.ent.interpretation_models import EntMessageInterpretation
app = create_app()
with app.app_context():
# Récupérer le contexte GMAO (partagé avec Outlook)
gmao_context = GmaoContext.query.first()
# Vérifier si activé
if not gmao_context or not gmao_context.auto_interpret_enabled:
log("Watchdog désactivé dans la config")
return False # Arrêter le watchdog
if not should_run(gmao_context):
# Pas le moment, attendre
return True
interval = gmao_context.auto_interpret_interval or 10
log(f"Exécution de l'interprétation (intervalle: {interval} min)...")
try:
# Synchroniser les messages
log("Démarrage synchronisation...")
sync_result = sync_ent_messages()
log(f"Sync result: {sync_result}")
# Utiliser le même contexte GMAO
max_messages = gmao_context.auto_interpret_max_emails or 5
lookback_hours = gmao_context.auto_interpret_lookback_hours or 24
cutoff_date = datetime.utcnow() - timedelta(hours=lookback_hours)
# Messages non interprétés
pending_messages = EntMessage.query.filter(
EntMessage.date >= cutoff_date
).filter(
~EntMessage.id.in_(db.session.query(EntMessageInterpretation.message_id))
).order_by(EntMessage.date.desc()).limit(max_messages).all()
processed_count = 0
created_count = 0
for msg in pending_messages:
try:
# Vérifier si l'expéditeur est exclu
excluded_senders = []
if gmao_context.excluded_senders:
try:
excluded_senders = json.loads(gmao_context.excluded_senders)
except json.JSONDecodeError:
excluded_senders = [gmao_context.excluded_senders]
if msg.sender_name and any(excl.lower() in msg.sender_name.lower() for excl in excluded_senders):
# Créer une interprétation rejetée
interpretation = EntMessageInterpretation(
message_id=msg.id,
analysis_type='information',
analysis_action='ignore',
status='rejected',
suggested_description=f"Message automatiquement ignoré car l'expéditeur '{msg.sender_name}' est dans la liste d'exclusion."
)
db.session.add(interpretation)
processed_count += 1
continue
# Analyser le message
result = analyze_ent_message_with_hermes(msg, gmao_context)
if result and not result.get('error'):
interpretation = EntMessageInterpretation(
message_id=msg.id,
analysis_type=result.get('type', 'information'),
analysis_priority=result.get('urgency', 'medium'),
analysis_summary=result.get('title', ''),
analysis_action=result.get('action', 'create'),
analysis_confidence=result.get('confidence', 0.5),
suggested_equipment=result.get('equipment'),
suggested_description=result.get('description'),
status='pending'
)
db.session.add(interpretation)
created_count += 1
elif result and result.get('error'):
log(f"Erreur analyse message {msg.id}: {result.get('error')}", level='error')
processed_count += 1
except Exception as e:
import traceback
log(f"Erreur interprétation message {msg.id}: {str(e)}", level='error')
log(f"Traceback: {traceback.format_exc()[-500:]}", level='error')
db.session.commit()
log(f"{processed_count} messages traités, {created_count} interprétations créées")
save_last_run()
log("save_last_run done")
db.session.commit()
log("db.session.commit done")
except Exception as e:
db.session.rollback()
import traceback
log(f"ERREUR dans le cycle: {str(e)}", level='error')
log(f"Traceback: {traceback.format_exc()[:500]}", level='error')
# Continuer quand même
try:
db.session.rollback()
except:
pass
return True # Continuer
def main():
"""Point d'entrée principal - pattern DND avec redémarrage automatique."""
global running
log("="*50)
log("Watchdog ENT démarré (pattern DND)")
log("="*50)
# Capturer le signal d'arrêt
def signal_handler(signum, frame):
global running
running = False
log("Signal reçu, arrêt du watchdog...")
signal.signal(signal.SIGTERM, signal_handler)
signal.signal(signal.SIGINT, signal_handler)
while running:
result = None
try:
result = run_watchdog()
if result is False:
log("Watchdog désactivé - arrêt")
break # Sortir proprement si désactivé
except Exception as e:
import traceback
log(f"Crash: {str(e)}", level='error')
log(f"Traceback: {traceback.format_exc()[-500:]}", level='error')
# Redémarrer automatiquement après crash
# Attendre avant la prochaine vérification
# Si le watchdog est désactivé, on attend plus longtemps pour éviter les redémarrages incessants
wait_seconds = 300 if result is False else 30
for _ in range(wait_seconds):
if not running:
break
time.sleep(1)
log("Watchdog ENT arrêté")
if __name__ == '__main__':
main()