import Foundation import Combine /// WebSocket client for real-time encounter correlation and messaging. /// /// This is the "live wire" of the interaction moment. When two present users /// are near each other, the server pushes a `mutual` event so both clients /// can trigger the reveal in real time. Revealed messages also arrive here /// (E2E ciphertext only — the server never sees plaintext). /// /// The client is a shared, observable singleton-style service so that both /// `RevealFlow` (mutual waves) and `ChatViewModel` (messages) consume the /// same socket. It auto-reconnects with backoff. @MainActor final class RealtimeClient: ObservableObject { // MARK: - Published state @Published private(set) var isConnected = false // MARK: - Events (Combine publishers) /// Emits a remote anonymous token when the server reports a mutual wave. let mutualEncounter = PassthroughSubject() /// Emits an incoming encrypted message. let incomingMessage = PassthroughSubject() // MARK: - Private private let socketURL = URL(string: "wss://api.proximity.app/ws")! private var task: URLSessionWebSocketTask? private var authToken: String? private var reconnectAttempts = 0 private var reconnectTask: Task? private var isActive = false // MARK: - Lifecycle /// Connect using the current session token. Idempotent. func connect(token: String) { authToken = token isActive = true reconnectAttempts = 0 openSocket() } func disconnect() { isActive = false reconnectTask?.cancel() task?.cancel(with: .goingAway, reason: nil) task = nil isConnected = false } private func openSocket() { guard isActive else { return } task?.cancel() var request = URLRequest(url: socketURL) if let authToken { request.setValue("Bearer \(authToken)", forHTTPHeaderField: "Authorization") } let session = URLSession(configuration: .default) let task = session.webSocketTask(with: request) self.task = task task.resume() receiveLoop() } // MARK: - Receiving private func receiveLoop() { task?.receive { [weak self] result in guard let self else { return } switch result { case .success(let message): self.isConnected = true self.reconnectAttempts = 0 switch message { case .data(let data): self.handle(data) case .string(let string): self.handle(Data(string.utf8)) @unknown default: break } self.receiveLoop() case .failure: self.handleDisconnect() } } } private func handleDisconnect() { isConnected = false guard isActive else { return } // Exponential backoff reconnect. let delay = min(pow(2.0, Double(reconnectAttempts)), 30.0) reconnectAttempts += 1 reconnectTask = Task { [weak self] in try? await Task.sleep(nanoseconds: UInt64(delay * 1_000_000_000)) guard !Task.isCancelled else { return } self?.openSocket() } } // MARK: - Event parsing private func handle(_ data: Data) { struct Event: Decodable { let type: String let remoteToken: String? let message: ChatMessage? } guard let event = try? JSONDecoder().decode(Event.self, from: data) else { return } switch event.type { case "mutual": if let token = event.remoteToken { mutualEncounter.send(token) } case "message": if let message = event.message { incomingMessage.send(message) } default: break } } }