- Add SYNC_STRATEGY_ESPOCRM_BASED.md detailing the sync flows and status management. - Create utilities for sync operations in services/beteiligte_sync_utils.py, including locking, timestamp comparison, conflict resolution, and notification handling. - Implement entity mapping between EspoCRM and Advoware in services/espocrm_mapper.py. - Develop a cron job for periodic sync checks in steps/vmh/beteiligte_sync_cron_step.py, emitting events for entities needing synchronization.
342 lines
12 KiB
Python
342 lines
12 KiB
Python
"""
|
|
Beteiligte Sync Utilities
|
|
|
|
Hilfsfunktionen für Sync-Operationen:
|
|
- Locking via syncStatus
|
|
- Timestamp-Vergleich
|
|
- Konfliktauflösung (EspoCRM wins)
|
|
- EspoCRM In-App Notifications
|
|
- Soft-Delete Handling
|
|
"""
|
|
|
|
from typing import Dict, Any, Optional, Tuple, Literal
|
|
from datetime import datetime
|
|
import pytz
|
|
import logging
|
|
from services.espocrm import EspoCRMAPI
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Timestamp-Vergleich Ergebnis-Typen
|
|
TimestampResult = Literal["espocrm_newer", "advoware_newer", "conflict", "no_change"]
|
|
|
|
|
|
class BeteiligteSync:
|
|
"""Utility-Klasse für Beteiligte-Synchronisation"""
|
|
|
|
def __init__(self, espocrm_api: EspoCRMAPI, context=None):
|
|
self.espocrm = espocrm_api
|
|
self.context = context
|
|
|
|
def _log(self, message: str, level: str = 'info'):
|
|
"""Logging mit Context-Support"""
|
|
if self.context and hasattr(self.context, 'logger'):
|
|
getattr(self.context.logger, level)(message)
|
|
else:
|
|
getattr(logger, level)(message)
|
|
|
|
async def acquire_sync_lock(self, entity_id: str) -> bool:
|
|
"""
|
|
Setzt syncStatus auf "syncing" (atomares Lock)
|
|
|
|
Args:
|
|
entity_id: EspoCRM CBeteiligte ID
|
|
|
|
Returns:
|
|
True wenn Lock erfolgreich, False wenn bereits im Sync
|
|
"""
|
|
try:
|
|
entity = await self.espocrm.get_entity('CBeteiligte', entity_id)
|
|
|
|
current_status = entity.get('syncStatus')
|
|
|
|
if current_status == 'syncing':
|
|
self._log(f"Entity {entity_id} bereits im Sync-Prozess", level='warning')
|
|
return False
|
|
|
|
# Setze Lock
|
|
await self.espocrm.update_entity('CBeteiligte', entity_id, {
|
|
'syncStatus': 'syncing'
|
|
})
|
|
|
|
self._log(f"Sync-Lock für {entity_id} erworben (vorher: {current_status})")
|
|
return True
|
|
|
|
except Exception as e:
|
|
self._log(f"Fehler beim Acquire Lock: {e}", level='error')
|
|
return False
|
|
|
|
async def release_sync_lock(
|
|
self,
|
|
entity_id: str,
|
|
new_status: str = 'clean',
|
|
error_message: Optional[str] = None,
|
|
increment_retry: bool = False
|
|
) -> None:
|
|
"""
|
|
Gibt Sync-Lock frei und setzt finalen Status
|
|
|
|
Args:
|
|
entity_id: EspoCRM CBeteiligte ID
|
|
new_status: Neuer syncStatus (clean, failed, conflict, etc.)
|
|
error_message: Optional: Fehlermeldung für syncErrorMessage
|
|
increment_retry: Ob syncRetryCount erhöht werden soll
|
|
"""
|
|
try:
|
|
update_data = {
|
|
'syncStatus': new_status,
|
|
'advowareLastSync': datetime.now(pytz.UTC).isoformat()
|
|
}
|
|
|
|
if error_message:
|
|
update_data['syncErrorMessage'] = error_message[:2000] # Max. 2000 chars
|
|
else:
|
|
update_data['syncErrorMessage'] = None
|
|
|
|
if increment_retry:
|
|
# Hole aktuellen Retry-Count
|
|
entity = await self.espocrm.get_entity('CBeteiligte', entity_id)
|
|
current_retry = entity.get('syncRetryCount') or 0
|
|
update_data['syncRetryCount'] = current_retry + 1
|
|
else:
|
|
update_data['syncRetryCount'] = 0
|
|
|
|
await self.espocrm.update_entity('CBeteiligte', entity_id, update_data)
|
|
|
|
self._log(f"Sync-Lock released: {entity_id} → {new_status}")
|
|
|
|
except Exception as e:
|
|
self._log(f"Fehler beim Release Lock: {e}", level='error')
|
|
|
|
@staticmethod
|
|
def parse_timestamp(ts: Any) -> Optional[datetime]:
|
|
"""
|
|
Parse verschiedene Timestamp-Formate zu datetime
|
|
|
|
Args:
|
|
ts: String, datetime oder None
|
|
|
|
Returns:
|
|
datetime-Objekt oder None
|
|
"""
|
|
if not ts:
|
|
return None
|
|
|
|
if isinstance(ts, datetime):
|
|
return ts
|
|
|
|
if isinstance(ts, str):
|
|
# EspoCRM Format: "2026-02-07 14:30:00"
|
|
# Advoware Format: "2026-02-07T14:30:00" oder "2026-02-07T14:30:00Z"
|
|
try:
|
|
# Entferne trailing Z falls vorhanden
|
|
ts = ts.rstrip('Z')
|
|
|
|
# Versuche verschiedene Formate
|
|
for fmt in [
|
|
'%Y-%m-%d %H:%M:%S',
|
|
'%Y-%m-%dT%H:%M:%S',
|
|
'%Y-%m-%d',
|
|
]:
|
|
try:
|
|
return datetime.strptime(ts, fmt)
|
|
except ValueError:
|
|
continue
|
|
|
|
# Fallback: ISO-Format
|
|
return datetime.fromisoformat(ts)
|
|
|
|
except Exception as e:
|
|
logger.warning(f"Konnte Timestamp nicht parsen: {ts} - {e}")
|
|
return None
|
|
|
|
return None
|
|
|
|
def compare_timestamps(
|
|
self,
|
|
espo_modified_at: Any,
|
|
advo_geaendert_am: Any,
|
|
last_sync_ts: Any
|
|
) -> TimestampResult:
|
|
"""
|
|
Vergleicht Timestamps und bestimmt Sync-Richtung
|
|
|
|
Args:
|
|
espo_modified_at: EspoCRM modifiedAt
|
|
advo_geaendert_am: Advoware geaendertAm
|
|
last_sync_ts: Letzter Sync (advowareLastSync)
|
|
|
|
Returns:
|
|
"espocrm_newer": EspoCRM wurde nach last_sync geändert und ist neuer
|
|
"advoware_newer": Advoware wurde nach last_sync geändert und ist neuer
|
|
"conflict": Beide wurden nach last_sync geändert
|
|
"no_change": Keine Änderungen seit last_sync
|
|
"""
|
|
espo_ts = self.parse_timestamp(espo_modified_at)
|
|
advo_ts = self.parse_timestamp(advo_geaendert_am)
|
|
sync_ts = self.parse_timestamp(last_sync_ts)
|
|
|
|
# Logging
|
|
self._log(
|
|
f"Timestamp-Vergleich: EspoCRM={espo_ts}, Advoware={advo_ts}, LastSync={sync_ts}",
|
|
level='debug'
|
|
)
|
|
|
|
# Falls kein last_sync → erster Sync, vergleiche direkt
|
|
if not sync_ts:
|
|
if not espo_ts or not advo_ts:
|
|
return "no_change"
|
|
|
|
if espo_ts > advo_ts:
|
|
return "espocrm_newer"
|
|
elif advo_ts > espo_ts:
|
|
return "advoware_newer"
|
|
else:
|
|
return "no_change"
|
|
|
|
# Check ob seit last_sync Änderungen
|
|
espo_changed = espo_ts and espo_ts > sync_ts
|
|
advo_changed = advo_ts and advo_ts > sync_ts
|
|
|
|
if espo_changed and advo_changed:
|
|
# Beide geändert seit last_sync → Konflikt
|
|
return "conflict"
|
|
elif espo_changed:
|
|
# Nur EspoCRM geändert
|
|
return "espocrm_newer" if (not advo_ts or espo_ts > advo_ts) else "conflict"
|
|
elif advo_changed:
|
|
# Nur Advoware geändert
|
|
return "advoware_newer"
|
|
else:
|
|
# Keine Änderungen
|
|
return "no_change"
|
|
|
|
async def send_notification(
|
|
self,
|
|
entity_id: str,
|
|
notification_type: Literal["conflict", "deleted"],
|
|
extra_data: Optional[Dict[str, Any]] = None
|
|
) -> None:
|
|
"""
|
|
Sendet EspoCRM In-App Notification
|
|
|
|
Args:
|
|
entity_id: CBeteiligte Entity ID
|
|
notification_type: "conflict" oder "deleted"
|
|
extra_data: Zusätzliche Daten für Nachricht
|
|
"""
|
|
try:
|
|
# Hole Entity-Daten
|
|
entity = await self.espocrm.get_entity('CBeteiligte', entity_id)
|
|
name = entity.get('name', 'Unbekannt')
|
|
betnr = entity.get('betnr')
|
|
assigned_user = entity.get('assignedUserId')
|
|
|
|
# Erstelle Nachricht basierend auf Typ
|
|
if notification_type == "conflict":
|
|
message = (
|
|
f"⚠️ Sync-Konflikt bei Beteiligten '{name}' (betNr: {betnr}). "
|
|
f"EspoCRM hat Vorrang - Änderungen wurden nach Advoware übertragen. "
|
|
f"Bitte prüfen Sie die Details."
|
|
)
|
|
elif notification_type == "deleted":
|
|
deleted_at = entity.get('advowareDeletedAt', 'unbekannt')
|
|
message = (
|
|
f"🗑️ Beteiligter '{name}' (betNr: {betnr}) wurde in Advoware gelöscht "
|
|
f"(am {deleted_at}). Der Datensatz wurde in EspoCRM markiert, aber nicht gelöscht. "
|
|
f"Bitte prüfen Sie, ob dies beabsichtigt war."
|
|
)
|
|
else:
|
|
message = f"Benachrichtigung für Beteiligten '{name}'"
|
|
|
|
# Erstelle Notification in EspoCRM
|
|
notification_data = {
|
|
'type': 'message',
|
|
'message': message,
|
|
'relatedType': 'CBeteiligte',
|
|
'relatedId': entity_id,
|
|
}
|
|
|
|
# Wenn assigned user vorhanden, sende an diesen
|
|
if assigned_user:
|
|
notification_data['userId'] = assigned_user
|
|
|
|
# Sende via API
|
|
result = await self.espocrm.api_call(
|
|
'Notification',
|
|
method='POST',
|
|
data=notification_data
|
|
)
|
|
|
|
self._log(f"Notification gesendet für {entity_id}: {notification_type}")
|
|
|
|
except Exception as e:
|
|
self._log(f"Fehler beim Senden der Notification: {e}", level='error')
|
|
|
|
async def handle_advoware_deleted(
|
|
self,
|
|
entity_id: str,
|
|
error_details: str
|
|
) -> None:
|
|
"""
|
|
Behandelt Fall dass Beteiligter in Advoware gelöscht wurde (404)
|
|
|
|
Args:
|
|
entity_id: CBeteiligte Entity ID
|
|
error_details: Fehlerdetails von Advoware API
|
|
"""
|
|
try:
|
|
now = datetime.now(pytz.UTC).isoformat()
|
|
|
|
# Update Entity: Soft-Delete Flag
|
|
await self.espocrm.update_entity('CBeteiligte', entity_id, {
|
|
'syncStatus': 'deleted_in_advoware',
|
|
'advowareDeletedAt': now,
|
|
'syncErrorMessage': f"Beteiligter existiert nicht mehr in Advoware. {error_details}"
|
|
})
|
|
|
|
self._log(f"Entity {entity_id} als deleted_in_advoware markiert")
|
|
|
|
# Sende Notification
|
|
await self.send_notification(entity_id, 'deleted')
|
|
|
|
except Exception as e:
|
|
self._log(f"Fehler beim Handle Deleted: {e}", level='error')
|
|
|
|
async def resolve_conflict_espocrm_wins(
|
|
self,
|
|
entity_id: str,
|
|
espo_entity: Dict[str, Any],
|
|
advo_entity: Dict[str, Any],
|
|
conflict_details: str
|
|
) -> None:
|
|
"""
|
|
Löst Konflikt auf: EspoCRM wins (überschreibt Advoware)
|
|
|
|
Args:
|
|
entity_id: CBeteiligte Entity ID
|
|
espo_entity: EspoCRM Entity-Daten
|
|
advo_entity: Advoware Entity-Daten
|
|
conflict_details: Details zum Konflikt
|
|
"""
|
|
try:
|
|
now = datetime.now(pytz.UTC).isoformat()
|
|
|
|
# Markiere als gelöst mit Konflikt-Info
|
|
await self.espocrm.update_entity('CBeteiligte', entity_id, {
|
|
'syncStatus': 'clean', # Gelöst!
|
|
'advowareLastSync': now,
|
|
'syncErrorMessage': f"Konflikt am {now}: {conflict_details}. EspoCRM hat gewonnen.",
|
|
'syncRetryCount': 0
|
|
})
|
|
|
|
self._log(f"Konflikt gelöst für {entity_id}: EspoCRM wins")
|
|
|
|
# Sende Notification
|
|
await self.send_notification(entity_id, 'conflict', {
|
|
'details': conflict_details
|
|
})
|
|
|
|
except Exception as e:
|
|
self._log(f"Fehler beim Resolve Conflict: {e}", level='error')
|