mana-swift-event-sync/Sources/ManaEventSync/Transport/SyncWSClient.swift
Till JS 423ba7c662 fix(event-sync): Crypto-Degradierung sichtbar machen + WS-Reconnect-Jitter
bootstrapCrypto: Bei aktivierter Verschlüsselung und unerreichbarem Vault
fiel die Engine still auf NoOpCryptoProvider zurück (nur Log.notice) —
Payloads wären unverschlüsselt rausgegangen, ohne dass die App es merkt.
Jetzt Log.error + onError(error), damit die App reagieren kann (z.B. Push
pausieren). Lokaler Betrieb läuft via NoOp weiter (Local-First).

SyncWSClient.scheduleReconnect: Equal Jitter (50-100%) auf den 1-30s-Backoff,
verhindert synchrones Reconnecten aller Clients nach flächigem Server-Ausfall
(Thundering Herd). Backoff-Verlauf bleibt deterministisch.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-02 14:42:44 +02:00

181 lines
6.2 KiB
Swift

import Foundation
import ManaCore
/// WebSocket-Client für `wss://<sync>/ws`. Optional — nur Apps, die
/// Live-Push wollen (sonst reicht das Poll-Intervall).
///
/// Wire-Protocol gespiegelt aus `mana/services/mana-sync/internal/ws/`
/// + `@mana/event-sync` `client/ws.ts`. Reconnect mit exponential
/// backoff (1s..30s), Reset bei `auth-ok`.
@MainActor
final class SyncWSClient {
private let wsURL: URL
private let appId: String
private let auth: AuthClient
private let onSyncAvailable: @MainActor (_ appId: String, _ aggregateId: String) -> Void
private var task: URLSessionWebSocketTask?
private var session: URLSession?
private var pingTimer: Task<Void, Never>?
private var reconnectTimer: Task<Void, Never>?
private var receiveTask: Task<Void, Never>?
private var stopped = false
private var currentBackoff: TimeInterval = 1
init(
syncURL: URL,
appId: String,
auth: AuthClient,
onSyncAvailable: @MainActor @escaping (String, String) -> Void
) {
var comp = URLComponents(url: syncURL, resolvingAgainstBaseURL: false) ?? URLComponents()
comp.scheme = comp.scheme == "http" ? "ws" : "wss"
comp.path = (comp.path.isEmpty || comp.path == "/") ? "/ws" : comp.path + "/ws"
wsURL = comp.url ?? syncURL
self.appId = appId
self.auth = auth
self.onSyncAvailable = onSyncAvailable
}
func start() async {
stopped = false
await connect()
}
func stop() {
stopped = true
pingTimer?.cancel()
pingTimer = nil
reconnectTimer?.cancel()
reconnectTimer = nil
receiveTask?.cancel()
receiveTask = nil
task?.cancel(with: .normalClosure, reason: nil)
task = nil
}
private func connect() async {
guard !stopped else { return }
session = URLSession(configuration: .default)
let webSocket = session!.webSocketTask(with: wsURL)
task = webSocket
webSocket.resume()
await sendAuth()
startReceiveLoop()
}
private func sendAuth() async {
do {
let token = try await auth.freshAccessToken()
try await sendJSON(["type": "auth", "token": token])
} catch {
Log.sync.error("ws auth send failed: \(error.localizedDescription, privacy: .public)")
scheduleReconnect()
}
}
private func sendSubscribe() async {
try? await sendJSON(["type": "subscribe", "appIds": [appId]])
}
private func sendPing() async {
try? await sendJSON(["type": "ping"])
}
private func sendJSON(_ payload: [String: Any]) async throws {
let data = try JSONSerialization.data(withJSONObject: payload, options: [])
guard let str = String(data: data, encoding: .utf8) else { return }
try await task?.send(.string(str))
}
private func startReceiveLoop() {
receiveTask?.cancel()
receiveTask = Task { [weak self] in
guard let self else { return }
while !Task.isCancelled {
guard let task else { break }
do {
let message = try await task.receive()
await handle(message)
} catch {
if !stopped {
Log.sync.error("ws receive failed: \(error.localizedDescription, privacy: .public)")
scheduleReconnect()
}
break
}
}
}
}
// swiftlint:disable:next cyclomatic_complexity
private func handle(_ message: URLSessionWebSocketTask.Message) async {
let text: String
switch message {
case let .string(value): text = value
case let .data(value): text = String(data: value, encoding: .utf8) ?? ""
@unknown default: return
}
guard
let data = text.data(using: .utf8),
let dict = try? JSONSerialization.jsonObject(with: data) as? [String: Any],
let type = dict["type"] as? String
else { return }
switch type {
case "auth-ok":
currentBackoff = 1
await sendSubscribe()
startPingTimer()
case "auth-failed":
let msg = (dict["humanMessage"] as? String) ?? "auth-failed"
Log.sync.error("ws auth failed: \(msg, privacy: .public)")
stop()
case "sync-available":
if let app = dict["appId"] as? String, let aggId = dict["aggregateId"] as? String {
onSyncAvailable(app, aggId)
}
case "subscribe-ok", "pong":
break
case "error":
let msg = (dict["humanMessage"] as? String) ?? "ws-error"
Log.sync.error("ws error: \(msg, privacy: .public)")
default:
break // forward-compat
}
}
private func startPingTimer() {
pingTimer?.cancel()
pingTimer = Task { [weak self] in
while !Task.isCancelled {
try? await Task.sleep(nanoseconds: 30_000_000_000)
if Task.isCancelled { break }
guard let self else { break }
await sendPing()
}
}
}
private func scheduleReconnect() {
guard !stopped, reconnectTimer == nil else { return }
pingTimer?.cancel()
pingTimer = nil
task?.cancel(with: .abnormalClosure, reason: nil)
task = nil
let capped = min(currentBackoff, 30)
currentBackoff = min(currentBackoff * 2, 30)
// Equal Jitter: tatsächlicher Delay zwischen 50% und 100% des Backoffs.
// Verhindert, dass nach einem flächigen Server-Ausfall alle Clients synchron
// zur selben Sekunde reconnecten (Thundering Herd). Der Backoff-Verlauf
// (1→2→4…30s) bleibt deterministisch, nur die Feuerzeit streut.
let delay = capped * Double.random(in: 0.5 ... 1.0)
reconnectTimer = Task { [weak self] in
try? await Task.sleep(nanoseconds: UInt64(delay * 1_000_000_000))
if Task.isCancelled { return }
guard let self else { return }
reconnectTimer = nil
await connect()
}
}
}