mana-swift-event-sync/Sources/ManaEventSync/Engine/EventSyncEngine.swift
Till JS 09963498fd fix: unused try?-Resultat in setSyncConsent (Warnung)
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 16:35:41 +02:00

702 lines
30 KiB
Swift

import Foundation
import ManaCore
import SwiftData
// Zentrale Engine — bewusst groß (Identität + Storage + Transport + Crypto).
// swiftlint:disable type_body_length file_length
/// Die app-facing Sync-Engine. Besitzt die Identität (anonym oder echt),
/// den lokalen Event-Log + Outbox, den Crypto-Hook und die Transport-
/// Orchestrierung (Pull/Push/Poll, optional WebSocket).
///
/// Lifecycle:
/// - `init(config:auth:)` — Storage + Crypto-Bootstrap
/// - `start()` — bei Login: Pull + Drain + Poll (+ WS); anonym: nur lokal
/// - `makeEnvelope(...)` + `emit(_:)` — Event schreiben (Identitäts-Stempel)
/// - `signIn(realUserId:)` — anonyme Events claimen + Sync starten
///
/// **Anonymous-First:** Ohne Login emittet die Engine lokal,
/// `attributedToUserId = anon:<id>`, die Outbox staut. `signIn` re-taggt
/// idempotent und drained. Der Server akzeptiert nur `attributedToUserId
/// == JWT.sub`, daher verlassen anonyme Events das Gerät nie vor dem Claim.
@MainActor
public final class EventSyncEngine {
private let config: EventSyncConfig
private let auth: AuthClient
private let storage: EventStorage
private let http: SyncHTTPClient
private var crypto: CryptoProvider = NoOpCryptoProvider()
private var ws: SyncWSClient?
private var pollTask: Task<Void, Never>?
private var pushInFlight = false
private var pullInFlight = false
private var cachedAnonId: String?
private var listeners: [UUID: @MainActor (String) -> Void] = [:]
/// Optionaler Fehler-Hook für die App-UI (z.B. „Sync fehlgeschlagen").
/// Wird bei Pull-/Push-Fehlern aufgerufen — zusätzlich zum Logging.
public var onError: (@MainActor (Error) -> Void)?
/// Dedizierter Hook für `422 schemaOutdated` (App-Schema älter als der
/// Server-Stand). Wird **statt** `onError` gefeuert und stoppt den Sync —
/// die App soll ein „Update verfügbar"-Banner zeigen, statt still
/// weiter gegen 422 zu pollen. Setzt voraus, dass die App einen
/// `schemaHash` pinnt (sonst sendet der Client keinen Header → kein 422).
public var onSchemaOutdated: (@MainActor () -> Void)?
/// Wird pro nicht-entschlüsselbarem Event aus dem Pull gefeuert. Das
/// Event wird übersprungen; der Verlust wird unabhängig vom Hook
/// persistent gezählt (``stats()`` → `decryptFailures`).
public var onDecryptFailure: (@MainActor (DecryptFailure) -> Void)?
/// Dedizierter Hook für `402 subscriptionRequired` (kein aktives
/// Cloud-Sync-Abo). Wird **statt** `onError` gefeuert; die App kann ein
/// „Cloud-Sync aktivieren"-Banner zeigen. Events bleiben in der Outbox
/// (kein Verlust), Auto-Recovery via Poll, sobald das Abo aktiv ist.
public var onSyncRequired: (@MainActor () -> Void)?
/// 402: Server blockiert Sync mangels Abo. Unterdrückt event-getriggerte
/// Push/Pull (Outbox staut sich lokal); der Poll prüft weiter (force) und
/// löst das Flag beim nächsten erfolgreichen Call.
private var syncRequired = false
/// `true`, wenn der Server Sync mangels Abo (402) blockiert.
public var isSyncRequired: Bool { syncRequired }
/// Cloud-Sync-Einwilligung (nur relevant, wenn `config.requireSyncConsent`).
/// Ohne Gate immer als erteilt behandelt. Beim `start()` aus mana-auth
/// geladen, per ``setSyncConsent(_:)`` umgeschaltet.
private var consentGranted = false
/// Effektive Einwilligung: ohne Gate (`requireSyncConsent == false`) immer
/// `true` (Verhalten wie bisher), sonst die geladene/gesetzte Nutzer-Wahl.
private var syncConsented: Bool { !config.requireSyncConsent || consentGranted }
/// `true`, wenn diese App das Opt-in-Gate nutzt (für die Konto-UI: nur dann
/// einen Sync-Toggle zeigen).
public var syncConsentRequired: Bool { config.requireSyncConsent }
/// Aktueller Einwilligungs-Stand (Initialwert für den Konto-Toggle).
public var hasSyncConsent: Bool { consentGranted }
/// Der SwiftData-Container des Event-Stores. Für Apps, die via
/// `EventSyncConfig.additionalModels` eigene `@Model`-Typen im selben
/// Store ablegen (z.B. nutriphis `LocalBMRProfile`) — die brauchen
/// einen `ModelContext` auf genau diesen Container.
public var modelContainer: ModelContainer { storage.container }
public init(config: EventSyncConfig, auth: AuthClient) throws {
self.config = config
self.auth = auth
storage = try EventStorage(
storeName: config.storeName,
appGroupIdentifier: config.appGroupIdentifier,
additionalModels: config.additionalModels,
inMemory: config.inMemory
)
http = SyncHTTPClient(
baseURL: config.syncURL,
appId: config.appId,
schemaHash: config.schemaHash,
auth: auth
)
}
// MARK: - Lifecycle
public func start() async {
// Opt-in-Gate: Einwilligung VOR der isSignedIn-Prüfung laden, da diese
// (über authMode) jetzt von der Einwilligung abhängt. Ohne Zustimmung
// bleibt die Engine lokal-only, auch bei vorhandenem Login.
if config.requireSyncConsent {
consentGranted = await loadInitialConsent()
}
await bootstrapCrypto()
guard isSignedIn else {
// Anonym (oder eingeloggt ohne Einwilligung): rein lokal, kein
// Netzwerk. Nur Retention pflegen.
enforceRetention()
return
}
await pull()
await drainOutbox()
startPollTimer()
if config.enableWebSocket { startWebSocket() }
}
public func stop() {
pollTask?.cancel()
pollTask = nil
ws?.stop()
ws = nil
}
// MARK: - Identität
/// Roher Auth-Modus aus dem ``AuthClient`` — ignoriert die Sync-
/// Einwilligung. Intern für die Consent-Logik (Initial-Load,
/// `setSyncConsent`), die wissen muss, ob überhaupt ein Login vorliegt.
private func rawAuthMode() -> AuthMode {
if case .signedIn = auth.status,
let token = try? auth.currentAccessToken(),
let sub = JWTSubject.subject(of: token)
{
return .signedIn(userId: sub)
}
return .anonymous(anonId: anonymousUserId())
}
public func authMode() -> AuthMode {
// Ohne Cloud-Sync-Einwilligung verhält sich die Engine lokal-only
// (anonym), auch wenn ein JWT vorliegt: kein Server-Egress, Events
// stauen in der Outbox und werden bei `setSyncConsent(true)` (bzw.
// `signIn`) in den Account geclaimt. Mirror zu Webs gegateter
// `tokenForSync` → `null`.
guard syncConsented else { return .anonymous(anonId: anonymousUserId()) }
return rawAuthMode()
}
// MARK: - Cloud-Sync-Einwilligung (Opt-in-Gate)
/// Schaltet die Cloud-Sync-Einwilligung zur Laufzeit um (Konto-Toggle).
/// No-op, wenn die App das Gate nicht nutzt (`requireSyncConsent == false`)
/// oder sich der Wert nicht ändert.
///
/// - `true` bei vorhandenem Login: claimt die bisher lokal/anonym
/// gestauten Events in den Account und startet den Sync (entspricht Webs
/// `applySyncConsent(true)` → `signIn`).
/// - `false`: stoppt den Server-Sync (Timer/WS), lokale Daten bleiben; neue
/// Events laufen wieder lokal-only. Bereits auf dem Server liegende Daten
/// werden **nicht** automatisch gelöscht (das ist der separate
/// VVT-§33-Purge, bewusst nicht hier verdrahtet).
public func setSyncConsent(_ granted: Bool) async {
guard config.requireSyncConsent, consentGranted != granted else { return }
consentGranted = granted
if granted {
// Consent ist hier bereits gesetzt (Caller hat ihn nach mana-auth
// persistiert) → direkt claimen, NICHT über `signIn` (das würde aus
// mana-auth neu laden und ggf. noch das alte `false` sehen).
if case let .signedIn(userId) = rawAuthMode() {
_ = try? await claimAndStartSync(realUserId: userId)
}
} else {
stop()
}
}
/// Liest die Initial-Einwilligung aus mana-auth `GET /api/v1/settings`
/// (App-Override vor globaler Präferenz, fehlend = aus). Ohne Login oder bei
/// jedem Fehler: `false` (local-first-Default).
private func loadInitialConsent() async -> Bool {
guard case .signedIn = rawAuthMode(), let token = try? auth.currentAccessToken() else {
return false
}
var request = URLRequest(url: config.authBaseURL.appendingPathComponent("api/v1/settings"))
request.setValue("Bearer \(token)", forHTTPHeaderField: "Authorization")
do {
let (data, response) = try await URLSession.shared.data(for: request)
guard let http = response as? HTTPURLResponse, http.statusCode == 200 else { return false }
return parseSyncConsent(fromSettingsJSON: data, appId: config.appId)
} catch {
return false
}
}
/// Stabile anonyme Identität (`anon:<guestId>`). Reused die
/// ManaCore-Guest-ID (App-Group-shareable), persistiert in Meta.
public func anonymousUserId() -> String {
if let cachedAnonId { return cachedAnonId }
if let value = try? storage.metaGet(EventStorageKey.anonymousUserId), !value.isEmpty {
cachedAnonId = value
return value
}
let base = auth.currentGuestId() ?? ((try? auth.enterGuestMode()) ?? UUID().uuidString)
let id = "anon:\(base)"
try? storage.metaSet(EventStorageKey.anonymousUserId, id)
cachedAnonId = id
return id
}
private var isSignedIn: Bool {
if case .signedIn = authMode() { return true }
return false
}
/// True, wenn ein gespeicherter anonymer Identitäts-Tag mit Events
/// existiert, der beim Login geclaimt werden müsste. Liest nur Meta +
/// Index — generiert (anders als ``anonymousUserId()``) **keine** neue
/// Anon-ID. Damit kann der Caller `signIn` bei jedem Login aufrufen,
/// ohne auf bereits eingeloggten Launches doppelt zu pullen.
public func hasPendingAnonymousEvents() -> Bool {
guard let anon = try? storage.metaGet(EventStorageKey.anonymousUserId),
!anon.isEmpty
else { return false }
return (try? storage.hasEvents(attributedTo: anon)) ?? false
}
// MARK: - Envelope-Bau (mit Identitäts-Stempel)
/// Baut ein Envelope und stempelt `attributedToUserId`/`actor` aus dem
/// aktuellen Auth-Modus. Ersetzt das frühere `currentUserId()` in den
/// App-Coordinators — der eine Hebel, der alles local-first macht.
public func makeEnvelope(
aggregateId: String,
eventType: String,
eventVersion: Int = 1,
payload: JSONValue,
causationId: String? = nil,
correlationId: String? = nil
) -> EventEnvelope {
let uid: String = switch authMode() {
case let .signedIn(userId): userId
case let .anonymous(anonId): anonId
}
let eventId = ULID.generate()
return EventEnvelope(
eventId: eventId,
aggregateId: aggregateId,
appId: config.appId,
eventType: eventType,
eventVersion: eventVersion,
occurredAt: nowISO(),
causationId: causationId,
correlationId: correlationId,
actor: .user(principalId: uid),
attributedToUserId: uid,
origin: "user",
idempotencyKey: eventId,
clientId: clientId,
payload: payload
)
}
// MARK: - Emit
/// Schreibt Event lokal (plaintext, für den Reducer) + Outbox
/// (Crypto-Wire) + triggert Push, wenn eingeloggt. Anonym: Outbox
/// staut, Retention greift.
/// Krypto-Kontext aus dem Envelope (für den ScopedCryptoProvider → scopeId
/// aus aggregateId; andere Provider ignorieren ihn).
private func cryptoContext(_ e: EventEnvelope) -> CryptoContext {
CryptoContext(aggregateId: e.aggregateId, attributedToUserId: e.attributedToUserId, appId: e.appId)
}
public func emit(_ envelope: EventEnvelope) async throws {
try storage.appendEvents([envelope])
let wirePayload = try await crypto.encryptPayload(envelope.payload, context: cryptoContext(envelope))
try storage.enqueueOutbox(envelope.withPayload(wirePayload))
notify(envelope.aggregateId)
if isSignedIn {
await drainOutbox()
} else {
enforceRetention()
}
}
// MARK: - Lesen (lokale Projektion)
public func events(forAggregate aggregateId: String) throws -> [EventEnvelope] {
try storage.eventsForAggregate(aggregateId)
}
public func events(forAggregatePrefix prefix: String) throws -> [String: [EventEnvelope]] {
try storage.eventsForAggregatePrefix(prefix)
}
// MARK: - Claim bei Login
/// Liftet alle anonymen Events in den Account: re-taggt EventLog +
/// Outbox idempotent von `anon:<id>` auf `realUserId`, bootstrappt
/// Crypto neu (Master-Key jetzt verfügbar) und startet den Sync.
/// `realUserId` ist der JWT-`sub` (Caller hat sich via AuthClient
/// eingeloggt). Stilles Re-Tagging, kein UI-Schritt (wie Web).
@discardableResult
public func signIn(realUserId: String) async throws -> ClaimResult {
// Adopting-Apps mit Opt-in-Gate: Login heißt nicht automatisch Sync. Die
// beim `start()` (vor diesem Login) geladene Einwilligung jetzt neu aus
// mana-auth ziehen; ohne Zustimmung wird NICHT geclaimt/gesynct (Events
// bleiben anonym lokal). Damit ist auch der bestehende
// `claimOnSignIn`→`signIn`-Pfad der Apps ohne Änderung consent-sicher.
if config.requireSyncConsent {
consentGranted = await loadInitialConsent()
guard consentGranted else {
return ClaimResult(rewrittenEvents: 0, newUserId: realUserId)
}
}
return try await claimAndStartSync(realUserId: realUserId)
}
/// Re-taggt anonyme Events auf den Account, bootstrappt Crypto, re-
/// verschlüsselt die Outbox und startet den Sync. Gemeinsamer Kern von
/// ``signIn(realUserId:)`` (nach dem Consent-Check) und
/// ``setSyncConsent(_:)`` (Toggle, Consent bereits gesetzt).
@discardableResult
private func claimAndStartSync(realUserId: String) async throws -> ClaimResult {
let anon = anonymousUserId()
let rewritten = try storage.reattribute(from: anon, to: realUserId)
cachedAnonId = nil
try? storage.metaSet(EventStorageKey.anonymousUserId, "")
await bootstrapCrypto()
await reencryptOutbox()
await pull()
await drainOutbox()
startPollTimer()
if config.enableWebSocket { startWebSocket() }
Log.sync
.notice(
"signIn claim: \(rewritten, privacy: .public) Event(s) auf \(realUserId, privacy: .public) re-getaggt"
)
return ClaimResult(rewrittenEvents: rewritten, newUserId: realUserId)
}
/// Stoppt den Sync (Timer/WS), lokale Daten bleiben. Das eigentliche
/// Abmelden macht der Caller über `AuthClient`.
public func signOut() {
stop()
}
/// Schreddert den Sub-Key eines Scopes (Crypto-Shredding, A3). Ohne
/// `shredAfter` sofort (privat, DSGVO Art. 17); mit → geplant (Firma, bis
/// zur GeBüV-Frist). Macht die synchronisierten/Cross-Device-Events des
/// Scopes permanent unentschlüsselbar. Verlangt Login.
public func shredScope(scopeId: String, shredAfter: Date? = nil) async throws {
let vault = ScopeVaultClient(auth: auth, authBaseURL: config.authBaseURL)
try await vault.shred(appId: config.appId, scopeId: scopeId, shredAfter: shredAfter)
}
/// Löscht lokal alle Events, deren aggregateId mit `idPrefix` beginnt
/// (A3-Crypto-Shredding-Begleiter: lokale Klartext-Residue eines gelöschten
/// Scopes wegräumen). Liefert die Anzahl gelöschter Events.
@discardableResult
public func purgeAggregatePrefix(_ idPrefix: String) throws -> Int {
try storage.purgeByAggregatePrefix(idPrefix)
}
/// Re-verschlüsselt die wartenden Outbox-Payloads mit dem aktuellen
/// Crypto-Provider. Beim Claim emittierte Guest-Events liegen plaintext
/// (NoOp); nach `bootstrapCrypto` (Master-Key da) sollen sie nicht
/// plaintext auf den Server — sonst bricht das „encrypted at rest" für
/// geclaimte Daten. No-op wenn Crypto NoOp bleibt.
private func reencryptOutbox() async {
guard let items = try? storage.pendingOutbox(limit: 100_000) else { return }
for item in items {
guard let envelope = try? item.envelope() else { continue }
do {
let ctx = cryptoContext(envelope)
let plain = try await crypto.decryptPayload(envelope.payload, context: ctx)
let rewrapped = try await crypto.encryptPayload(plain, context: ctx)
try storage.updateOutboxEnvelope(eventId: envelope.eventId, envelope.withPayload(rewrapped))
} catch {
Log.sync.error("re-encrypt outbox \(envelope.eventId, privacy: .public) failed")
}
}
}
// MARK: - Aggregate-Change-Notifications
@discardableResult
public func onAggregateChange(_ listener: @MainActor @escaping (String) -> Void) -> UUID {
let token = UUID()
listeners[token] = listener
return token
}
public func unsubscribe(_ token: UUID) {
listeners.removeValue(forKey: token)
}
// MARK: - Diagnose / Wartung
public func stats() throws -> EngineStats {
try EngineStats(
eventsLocal: storage.countEvents(),
outboxPending: storage.outboxCount(),
lastSyncAt: storage.metaGet(EventStorageKey.lastSyncAt),
authMode: authMode(),
decryptFailures: storage.metaGet(EventStorageKey.decryptFailureCount).flatMap(Int.init) ?? 0,
decryptFailureLastAt: storage.metaGet(EventStorageKey.decryptFailureLastAt)
)
}
/// Setzt den Decrypt-Failure-Zähler zurück — nach erfolgreichem
/// Recovery (Key wieder verfügbar + Re-Pull).
public func clearDecryptFailures() throws {
try storage.metaSet(EventStorageKey.decryptFailureCount, "")
try storage.metaSet(EventStorageKey.decryptFailureLastAt, "")
try storage.metaSet(EventStorageKey.decryptFailureLastEvent, "")
}
public func clearLocal() throws {
try storage.clearAll()
cachedAnonId = nil
}
// MARK: - Internals
private func notify(_ aggregateId: String) {
for listener in listeners.values {
listener(aggregateId)
}
}
private func nowISO() -> String {
ISO8601DateFormatter().string(from: .init())
}
private lazy var clientId: String = {
if let value = try? storage.metaGet(EventStorageKey.clientId), !value.isEmpty {
return value
}
let id = ULID.generate()
try? storage.metaSet(EventStorageKey.clientId, id)
return id
}()
private func bootstrapCrypto() async {
guard config.enableEncryption, isSignedIn else {
crypto = NoOpCryptoProvider()
return
}
do {
if let resolver = config.scopeResolver {
crypto = try await createScopedKeyProviderFromVault(
auth: auth,
authBaseURL: config.authBaseURL,
appId: config.appId,
scopeResolver: resolver
)
} else {
crypto = try await createMasterKeyProviderFromVault(auth: auth, authBaseURL: config.authBaseURL)
}
} catch {
crypto = NoOpCryptoProvider()
// Verschlüsselung war angefordert, der Vault ist aber nicht erreichbar.
// Wir degradieren NICHT stillschweigend (vorher nur `.notice`): Log auf
// `.error` + `onError`, damit die App weiß, dass Payloads bis zum nächsten
// erfolgreichen Bootstrap unverschlüsselt wären. Lokaler Betrieb läuft via
// NoOp weiter; ob Push deshalb pausiert werden soll, entscheidet die App
// im `onError`-Handler.
Log.sync.error("Vault unerreichbar trotz aktivierter Verschlüsselung — NoOp-Crypto-Fallback, Payloads unverschlüsselt: \(error.localizedDescription, privacy: .public)")
onError?(error)
}
}
private func enforceRetention() {
let retention = config.anonymousRetention
guard retention.maxEvents > 0 || retention.maxAgeDays > 0 else { return }
try? storage.pruneAnonymous(maxEvents: retention.maxEvents, maxAgeDays: retention.maxAgeDays)
}
private func startWebSocket() {
guard ws == nil else { return }
ws = SyncWSClient(
syncURL: config.syncURL,
appId: config.appId,
auth: auth,
onSyncAvailable: { [weak self] _, _ in
Task { await self?.pull() }
}
)
Task { await ws?.start() }
}
private func startPollTimer() {
pollTask?.cancel()
pollTask = Task { [weak self] in
guard let self else { return }
let interval = config.pollInterval
while !Task.isCancelled {
try? await Task.sleep(nanoseconds: UInt64(interval * 1_000_000_000))
if Task.isCancelled { break }
// force: true — der Poll prüft auch im sync-required-Zustand
// weiter und erholt sich automatisch, sobald das Abo aktiv ist.
await pull(force: true)
await drainOutbox(force: true)
}
}
}
private func pull(force: Bool = false) async {
guard isSignedIn, !pullInFlight else { return }
// Im 402-Zustand keine event-getriggerten Pulls — nur der Poll (force).
if syncRequired, !force { return }
pullInFlight = true
defer { pullInFlight = false }
do {
var cursor: Int64 = 0
if let stored = try storage.metaGet(EventStorageKey.lastPullCursor), let value = Int64(stored) {
cursor = value
}
while true {
let res = try await http.pull(since: cursor, limit: 100)
if res.events.isEmpty { break }
var decrypted: [EventEnvelope] = []
decrypted.reserveCapacity(res.events.count)
for event in res.events {
do {
let plain = try await decryptPulledPayload(event.payload, eventId: event.eventId, crypto: crypto, context: cryptoContext(event))
decrypted.append(event.withPayload(plain))
} catch {
recordDecryptFailure(event, error: error)
}
}
try storage.appendEvents(decrypted)
for aggId in Set(decrypted.map(\.aggregateId)) {
notify(aggId)
}
if let last = res.events.last, let seq = last.sequenceNumber {
cursor = seq
try storage.metaSet(EventStorageKey.lastPullCursor, String(seq))
}
if !res.hasMore { break }
}
try storage.metaSet(EventStorageKey.lastSyncAt, nowISO())
// Erfolgreicher Pull → etwaige 402-Blockade aufheben.
clearSyncRequired()
} catch {
// 402: kein aktives Abo. Cursor unverändert, kein Schaden;
// Auto-Recovery via Poll, sobald aktiv.
if isSubscriptionRequired(error) {
enterSyncRequired()
return
}
reportSyncError(error)
}
}
/// `true`, wenn Verschlüsselung angefordert ist, der Crypto-Provider aber
/// auf NoOp degradierte (Vault-Bootstrap fehlgeschlagen). In diesem Zustand
/// dürfen eingeloggte Account-Daten NICHT als Klartext gepusht werden.
/// (Apps mit `enableEncryption == false` nutzen NoOp bewusst → kein Block.)
private var isCryptoDegraded: Bool {
config.enableEncryption && crypto.providerId == NoOpCryptoProvider.id
}
private func drainOutbox(force: Bool = false) async {
guard isSignedIn, !pushInFlight else { return }
// Im 402-Zustand keine event-getriggerten Pushes — nur der Poll (force).
if syncRequired, !force { return }
// Schutz vor Klartext-Leak: Lieber in der Outbox stauen als unverschlüsselt
// rausschicken. Ein erneut erfolgreicher `bootstrapCrypto` (nächster
// `signIn`/`start`) re-wrapped via `reencryptOutbox` und flusht dann.
if isCryptoDegraded {
Log.sync.error("drainOutbox pausiert: Verschlüsselung aktiv, Crypto aber auf NoOp (Vault unerreichbar) — Events bleiben in der Outbox")
return
}
pushInFlight = true
defer { pushInFlight = false }
do {
while true {
let items = try storage.pendingOutbox(limit: 100)
if items.isEmpty { break }
let envelopes = try items.map { try $0.envelope() }
do {
let res = try await http.append(envelopes)
try storage.removeOutbox(eventIds: res.accepted.map(\.eventId))
if !res.rejected.isEmpty {
for rej in res.rejected {
try storage.incrementOutboxAttempt(
eventId: rej.eventId,
error: "\(rej.errorCode): \(rej.humanMessage ?? "")"
)
}
Log.sync.error("\(res.rejected.count) Event(s) von sync2 abgelehnt")
break
}
// Erfolgreicher Push ohne Rejects → 402-Blockade aufheben.
clearSyncRequired()
} catch {
// 422 = App zu alt: Events bleiben unangetastet in der
// Outbox (kein Attempt-Increment), Sync stoppt, Banner.
if isSchemaOutdated(error) {
reportSyncError(error)
break
}
// 402 = kein Abo: Events BLEIBEN in der Outbox, KEIN
// Attempt-Increment (kein Verlust, kein Backoff-Aufblähen).
if isSubscriptionRequired(error) {
enterSyncRequired()
return
}
for item in items {
try storage.incrementOutboxAttempt(eventId: item.eventId, error: error.localizedDescription)
}
reportSyncError(error)
break
}
}
} catch {
reportSyncError(error)
}
}
/// Ein nicht-entschlüsselbares Event darf weder den Pull abbrechen noch
/// still verschwinden: skip + persistenter Zähler + Hook. Vorher war das
/// ein reiner Log-Eintrag — die App sah nie, dass ihr Events fehlen.
private func recordDecryptFailure(_ event: EventEnvelope, error: Error) {
let current = (try? storage.metaGet(EventStorageKey.decryptFailureCount)).flatMap { $0 }
.flatMap(Int.init) ?? 0
try? storage.metaSet(EventStorageKey.decryptFailureCount, String(current + 1))
try? storage.metaSet(EventStorageKey.decryptFailureLastAt, nowISO())
try? storage.metaSet(EventStorageKey.decryptFailureLastEvent, event.eventId)
Log.sync.error(
"decrypt failed for \(event.eventId, privacy: .public) (\(event.eventType, privacy: .public)) — skipped, Zähler \(current + 1, privacy: .public)"
)
let failure = DecryptFailure(
eventId: event.eventId,
aggregateId: event.aggregateId,
eventType: event.eventType,
error: error
)
onDecryptFailure?(failure)
onError?(EventSyncError.decryptFailed(eventId: event.eventId))
}
// MARK: - Fehler-Klassifikation
private func isSchemaOutdated(_ error: Error) -> Bool {
if let syncError = error as? EventSyncError, case .schemaOutdated = syncError { return true }
return false
}
private func isSubscriptionRequired(_ error: Error) -> Bool {
if let syncError = error as? EventSyncError, case .subscriptionRequired = syncError { return true }
return false
}
/// 402 vom Server: kein aktives Cloud-Sync-Abo. Setzt das Flag (unterdrückt
/// event-getriggerte Push/Pull), feuert `onSyncRequired`. Events bleiben in
/// der Outbox; der Poll (force) erholt sich automatisch, sobald aktiv.
private func enterSyncRequired() {
if !syncRequired {
syncRequired = true
Log.sync.error("402 — kein aktives Cloud-Sync-Abo; Events bleiben lokal in der Outbox")
onSyncRequired?()
}
}
/// Hebt die 402-Blockade auf (nach erfolgreichem Push/Pull). Der nächste
/// forcierte Poll-Drain zieht die gestaute Outbox nach.
private func clearSyncRequired() {
guard syncRequired else { return }
syncRequired = false
}
/// Klassifiziert einen Sync-Fehler. `422 schemaOutdated` (App-Schema
/// älter als Server) **stoppt** den Sync — weiteres Pollen würde nur
/// weiter 422en — und feuert den dedizierten `onSchemaOutdated`-Hook
/// (App zeigt „Update verfügbar"). Alles andere geht an `onError`.
private func reportSyncError(_ error: Error) {
if isSchemaOutdated(error) {
Log.sync.error("schema-hash drift (422) — App zu alt, Sync gestoppt bis Update")
stop()
onSchemaOutdated?()
} else {
Log.sync.error("sync failed: \(error.localizedDescription, privacy: .public)")
onError?(error)
}
}
}
// swiftlint:enable type_body_length file_length