CardSync/backend/sync_engine.py

398 lines
16 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.

import logging
import asyncio
from datetime import datetime, timezone, timedelta
from sqlalchemy.orm import Session
from models import User, SyncLog, SynologyConfig
from carddav_client import CardDAVClient
from ms_graph import get_valid_token, graph_get_contacts, graph_create_contact, graph_update_contact, graph_delete_contact
logger = logging.getLogger(__name__)
# Trennzeichen für UID im personalNotes-Feld (um Synology-UID in MS zu speichern)
UID_MARKER = "[cardsync_uid:"
# ── CardDAV Cache ──────────────────────────────────────────────────────────
# Bei vielen Usern die dasselbe Adressbuch synchronisieren wird die Synology
# sonst pro Sync-Zyklus N-mal abgefragt. Cache reduziert das auf 1× pro TTL.
# Cache pro (server_url, folder_url) — pro Adressbuch separater Eintrag.
_carddav_cache: dict = {}
_carddav_cache_locks: dict = {}
CARDDAV_CACHE_TTL_SECONDS = 300 # 5 Minuten
def _cache_key(server_url: str, folder_url: str) -> str:
return f"{server_url}||{folder_url}"
async def _get_cached_carddav_contacts(server_url: str, username: str, password: str, folder_url: str) -> list[dict]:
"""Lädt CardDAV-Kontakte aus dem Cache oder frisch — pro (Server, Folder)
nur ein paralleler Abruf dank Lock. Folge-Requests warten und bekommen das Ergebnis."""
key = _cache_key(server_url, folder_url)
now = datetime.now(timezone.utc)
# Cache-Hit?
entry = _carddav_cache.get(key)
if entry:
expires_at = entry["expires_at"]
if now < expires_at:
logger.info(f"CardDAV Cache HIT für {folder_url} (age: {(now - entry['cached_at']).total_seconds():.0f}s)")
return entry["contacts"]
# Lock pro Schlüssel — verhindert dass 170 User gleichzeitig CardDAV abfragen
lock = _carddav_cache_locks.setdefault(key, asyncio.Lock())
async with lock:
# Nochmal prüfen — ein anderer Task hat eventuell schon geladen
entry = _carddav_cache.get(key)
if entry and now < entry["expires_at"]:
logger.info(f"CardDAV Cache HIT (nach Lock) für {folder_url}")
return entry["contacts"]
logger.info(f"CardDAV Cache MISS für {folder_url} — lade frisch")
client = CardDAVClient(server_url, username, password)
# client.get_contacts ist synchron — in Thread auslagern damit asyncio nicht blockiert
contacts = await asyncio.to_thread(client.get_contacts, folder_url)
_carddav_cache[key] = {
"contacts": contacts,
"cached_at": now,
"expires_at": now + timedelta(seconds=CARDDAV_CACHE_TTL_SECONDS),
}
return contacts
def invalidate_carddav_cache(folder_url: str = None):
"""Cache invalidieren — entweder komplett oder nur für ein Adressbuch."""
if folder_url is None:
_carddav_cache.clear()
logger.info("CardDAV-Cache komplett geleert")
else:
keys_to_remove = [k for k in _carddav_cache if k.endswith(f"||{folder_url}")]
for k in keys_to_remove:
del _carddav_cache[k]
logger.info(f"CardDAV-Cache invalidiert für {folder_url}")
def _embed_uid(uid: str, notes: str = "") -> str:
"""UID in personalNotes einbetten, damit wir MS-Kontakte mit vCards matchen können"""
if not uid:
return notes
marker = f"{UID_MARKER}{uid}]"
if marker in (notes or ""):
return notes
return f"{notes}\n{marker}".strip() if notes else marker
def _extract_uid(notes: str) -> str:
"""UID aus personalNotes extrahieren"""
if not notes or UID_MARKER not in notes:
return ""
start = notes.index(UID_MARKER) + len(UID_MARKER)
end = notes.index("]", start)
return notes[start:end]
def _clean_payload(contact: dict) -> dict:
"""Bereinigt das Kontakt-Dict für die Microsoft Graph API.
- Entfernt interne Felder (mit _ prefix)
- Entfernt leere Strings (Graph erwartet null oder weglassen)
- Entfernt leere Adress-Objekte (alle Felder leer -> komplettes Objekt weg)
- Stellt sicher dass Listen wirklich Listen mit Inhalt sind
"""
payload = {}
for k, v in contact.items():
if k.startswith("_"):
continue
# Leere Strings -> weg
if isinstance(v, str):
if v.strip():
payload[k] = v
continue
# Adress-Objekte: nur senden wenn min. 1 Feld gefüllt
if isinstance(v, dict):
cleaned = {sk: sv for sk, sv in v.items() if isinstance(sv, str) and sv.strip()}
if cleaned:
payload[k] = cleaned
continue
# Listen: leere raus
if isinstance(v, list):
if k == "emailAddresses":
# nur Einträge mit gültiger address behalten
cleaned_list = [e for e in v if isinstance(e, dict) and e.get("address", "").strip()]
# Graph-Schema: email braucht nur "address", "name" ist optional; "type" entfernen!
cleaned_list = [{"address": e["address"], "name": e.get("name", "") or e["address"]} for e in cleaned_list]
# Microsoft Graph erlaubt max. 3 E-Mail-Adressen
if len(cleaned_list) > 3:
cleaned_list = cleaned_list[:3]
if cleaned_list:
payload[k] = cleaned_list
elif k in ("businessPhones", "homePhones"):
cleaned_list = [p for p in v if isinstance(p, str) and p.strip()]
# Microsoft Graph erlaubt max. 2 Einträge pro Telefon-Liste
if len(cleaned_list) > 2:
cleaned_list = cleaned_list[:2]
if cleaned_list:
payload[k] = cleaned_list
else:
if v:
payload[k] = v
continue
# Sonst: nur setzen wenn Wert da
if v is not None and v != "":
payload[k] = v
return payload
def _embed_uid(uid: str, notes: str = "") -> str:
"""UID in personalNotes einbetten, damit wir MS-Kontakte mit vCards matchen können"""
if not uid:
return notes
marker = f"{UID_MARKER}{uid}]"
if marker in (notes or ""):
return notes
return f"{notes}\n{marker}".strip() if notes else marker
def _extract_uid(notes: str) -> str:
"""UID aus personalNotes extrahieren"""
if not notes or UID_MARKER not in notes:
return ""
start = notes.index(UID_MARKER) + len(UID_MARKER)
end = notes.index("]", start)
return notes[start:end]
def _normalize_addr(addr: dict) -> dict:
"""Normalisiert ein Adress-Dict für stabilen Vergleich"""
if not isinstance(addr, dict):
return {}
return {
"street": (addr.get("street") or "").strip(),
"city": (addr.get("city") or "").strip(),
"state": (addr.get("state") or "").strip(),
"postalCode": (addr.get("postalCode") or "").strip(),
"countryOrRegion": (addr.get("countryOrRegion") or "").strip(),
}
def _contacts_equal(carddav_contact: dict, ms_contact: dict) -> bool:
"""Prüft ob sich der Kontakt geändert hat.
Vergleicht ALLE Felder die wir auch hochladen — string-Felder, Listen, Adressen.
Wenn diese Funktion True liefert, wird der Kontakt nicht erneut hochgeladen.
"""
# Bereinigte Versionen vergleichen, damit Vergleich konsistent ist
cv_clean = _clean_payload(carddav_contact)
# MS-Kontakt: gleicher Cleanup für fairen Vergleich
# Notizen: cardsync-UID-Marker ausblenden, da der in MS bewusst gesetzt wird
def strip_uid_marker(notes: str) -> str:
if not notes:
return ""
if UID_MARKER not in notes:
return notes.strip()
# Marker entfernen
start = notes.index(UID_MARKER)
end = notes.index("]", start) + 1
return (notes[:start] + notes[end:]).strip()
# Einfache String-Felder
string_fields = [
"givenName", "surname", "middleName", "displayName", "title", "generation",
"jobTitle", "companyName", "department", "mobilePhone",
"businessHomePage", "nickName",
]
for f in string_fields:
if (cv_clean.get(f, "") or "") != (ms_contact.get(f, "") or ""):
return False
# Notizen vergleichen — UID-Marker rausrechnen
cv_notes = strip_uid_marker(cv_clean.get("personalNotes", ""))
ms_notes = strip_uid_marker(ms_contact.get("personalNotes", ""))
if cv_notes != ms_notes:
return False
# E-Mails
cv_emails = set((e.get("address") or "").lower() for e in cv_clean.get("emailAddresses", []))
ms_emails = set((e.get("address") or "").lower() for e in ms_contact.get("emailAddresses", []))
if cv_emails != ms_emails:
return False
# Telefone (geordnete Listen, da Graph nur max. 2 erlaubt sollte das robust sein)
if (cv_clean.get("businessPhones") or []) != (ms_contact.get("businessPhones") or []):
return False
if (cv_clean.get("homePhones") or []) != (ms_contact.get("homePhones") or []):
return False
# Adressen
if _normalize_addr(cv_clean.get("businessAddress", {})) != _normalize_addr(ms_contact.get("businessAddress", {})):
return False
if _normalize_addr(cv_clean.get("homeAddress", {})) != _normalize_addr(ms_contact.get("homeAddress", {})):
return False
return True
async def sync_user_contacts(user_id: int, db: Session) -> dict:
"""Synchronisiert Kontakte für einen Benutzer: Synology → MS365"""
user = db.query(User).filter(User.id == user_id).first()
if not user:
return {"status": "error", "message": "Benutzer nicht gefunden"}
if not user.sync_enabled:
return {"status": "skipped", "message": "Sync deaktiviert"}
if not user.carddav_folder:
return {"status": "error", "message": "Kein CardDAV-Ordner konfiguriert"}
# Sync-Log anlegen
log = SyncLog(user_id=user.id, started_at=datetime.now(timezone.utc))
db.add(log)
db.commit()
user.last_sync_status = "running"
user.updated_at = datetime.now(timezone.utc)
db.commit()
try:
# 1. Synology-Konfiguration laden
config = db.query(SynologyConfig).first()
if not config:
raise ValueError("Keine Synology-Konfiguration vorhanden")
# 2. CardDAV-Kontakte laden (mit Cache — bei vielen Usern mit gleichem
# Adressbuch wird Synology nur einmal pro Cache-TTL abgefragt)
carddav_contacts = await _get_cached_carddav_contacts(
config.server_url, config.username, config.password, user.carddav_folder
)
logger.info(f"[{user.email}] {len(carddav_contacts)} Kontakte von CardDAV geladen")
# 3. MS365-Kontakte laden
access_token = await get_valid_token(user, db)
if not access_token:
raise ValueError("Kein gültiges Microsoft-Token Benutzer muss sich neu anmelden")
# Bei Gruppen-Import-Usern: über /users/{upn}/contacts gehen (App-Token)
user_principal = user.email if getattr(user, "source", "oauth") == "group_import" else None
ms_contacts = await graph_get_contacts(access_token, user_principal)
logger.info(f"[{user.email}] {len(ms_contacts)} Kontakte von MS365 geladen")
# 4. Index aufbauen: UID → MS-Kontakt
ms_by_uid = {}
ms_by_display = {}
for ms_c in ms_contacts:
uid = _extract_uid(ms_c.get("personalNotes", ""))
if uid:
ms_by_uid[uid] = ms_c
name = ms_c.get("displayName", "").strip().lower()
if name:
ms_by_display[name] = ms_c
# 5. CardDAV → MS365 synchronisieren
carddav_uids = set()
created = updated = skipped = failed = 0
for contact in carddav_contacts:
uid = contact.get("_uid", "")
display = contact.get("displayName", "").strip().lower()
if not contact.get("displayName"):
continue # Kontakte ohne Namen überspringen
carddav_uids.add(uid)
# UID in Notes einbetten
notes = contact.get("personalNotes", "")
contact["personalNotes"] = _embed_uid(uid, notes)
# Felder für Graph API bereinigen
ms_payload = _clean_payload(contact)
if not ms_payload.get("displayName") and not ms_payload.get("givenName") and not ms_payload.get("surname"):
continue # Graph braucht mindestens einen Namen
try:
if uid and uid in ms_by_uid:
# Kontakt existiert schon → Update wenn nötig
ms_c = ms_by_uid[uid]
if not _contacts_equal(contact, ms_c):
await graph_update_contact(access_token, ms_c["id"], ms_payload, user_principal)
updated += 1
logger.debug(f"[{user.email}] Updated: {contact.get('displayName')}")
else:
skipped += 1
elif display and display in ms_by_display and not uid:
# Kein UID aber gleicher Name → Update + UID nachpflegen
ms_c = ms_by_display[display]
if not _contacts_equal(contact, ms_c):
await graph_update_contact(access_token, ms_c["id"], ms_payload, user_principal)
updated += 1
else:
skipped += 1
else:
# Neuer Kontakt → erstellen
await graph_create_contact(access_token, ms_payload, user_principal)
created += 1
logger.debug(f"[{user.email}] Created: {contact.get('displayName')}")
except Exception as ce:
failed += 1
logger.warning(f"[{user.email}] Kontakt '{contact.get('displayName')}' übersprungen: {ce}")
continue
# 6. Gelöschte Kontakte aus MS365 entfernen
deleted = 0
for uid, ms_c in ms_by_uid.items():
if uid not in carddav_uids:
try:
await graph_delete_contact(access_token, ms_c["id"], user_principal)
deleted += 1
logger.debug(f"[{user.email}] Deleted: {ms_c.get('displayName')}")
except Exception as de:
logger.warning(f"[{user.email}] Delete failed for {ms_c.get('displayName')}: {de}")
total = created + updated
message = f"{created} erstellt, {updated} aktualisiert, {deleted} gelöscht, {skipped} unverändert"
if failed:
message += f", {failed} fehlgeschlagen"
logger.info(f"[{user.email}] Sync abgeschlossen: {message}")
# Log abschließen
now = datetime.now(timezone.utc)
log.finished_at = now
log.status = "success"
log.contacts_synced = len(carddav_contacts)
log.contacts_created = created
log.contacts_updated = updated
log.contacts_deleted = deleted
user.last_sync_at = now
user.last_sync_status = "success"
user.last_sync_message = message
user.updated_at = now
db.commit()
return {"status": "success", "message": message, "created": created, "updated": updated, "deleted": deleted}
except Exception as e:
error_msg = str(e)
logger.error(f"[{user.email}] Sync failed: {error_msg}")
now = datetime.now(timezone.utc)
log.finished_at = now
log.status = "error"
log.error_message = error_msg
user.last_sync_at = now
user.last_sync_status = "error"
user.last_sync_message = error_msg
user.updated_at = now
db.commit()
return {"status": "error", "message": error_msg}