398 lines
16 KiB
Python
398 lines
16 KiB
Python
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}
|