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>
181 lines
6.2 KiB
Swift
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()
|
|
}
|
|
}
|
|
}
|