702 lines
30 KiB
Swift
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
|