// // SDLSuperClient.swift // Tun // // Created by 安礼成 on 2026/2/13. // import Foundation import Network final class SDLSuperClient: @unchecked Sendable { private let queue = DispatchQueue(label: "com.sdl.SuperClient.queue") // 专用队列保证线程安全 public let messageStream: AsyncThrowingStream private let messageContinuation: AsyncThrowingStream.Continuation private let readySignal = AsyncOneShot() private let stateLock = NSLock() private var isStarted = false private var isStopped = false private var isMessageContinuationFinished = false private let connection: NWConnection private let maxBufferSize: Int init(serverEndpoint: SDLConfiguration.ResolvedServerEndpoint, port: UInt16, maxBufferSize: Int = 2 * 1024 * 1024) { self.maxBufferSize = maxBufferSize let pairs = AsyncThrowingStream.makeStream(of: SDLSuperMessage.self, bufferingPolicy: .bufferingNewest(1024)) self.messageStream = pairs.stream self.messageContinuation = pairs.continuation let options = NWProtocolTLS.Options() serverEndpoint.host.withCString { sec_protocol_options_set_tls_server_name(options.securityProtocolOptions, $0) } sec_protocol_options_add_tls_application_protocol( options.securityProtocolOptions, "punchnet/1.0" ) // 这里设置证书的校验逻辑 sec_protocol_options_set_verify_block( options.securityProtocolOptions, { _, trust, complete in // 执行公钥校验 complete(SDLSuperTLSVerifier.verify(trust: trust, host: serverEndpoint.host)) }, queue ) let params = NWParameters(tls: options) // 关键:让 Network.framework 忽略系统代理 params.preferNoProxies = true self.connection = NWConnection(host: Self.makeEndpointHost(address: serverEndpoint.ip), port: .init(rawValue: port)!, using: params) SDLLogger.log("[SDLSuperClient] start with tls protocol", category: .super) } func run() async throws { guard self.markStarted() else { return } defer { self.stop() } do { try await withTaskCancellationHandler { self.connection.stateUpdateHandler = { [weak self] state in self?.handleConnectionStateUpdate(state) } self.connection.start(queue: self.queue) try await self.readySignal.wait() try await self.readLoop() self.finishMessageStream() } onCancel: { self.connection.stateUpdateHandler = nil self.connection.cancel() } } catch is CancellationError { self.finishMessageStream() throw CancellationError() } catch { self.finishMessageStream(throwing: error) throw error } } private static func makeEndpointHost(address ip: String) -> NWEndpoint.Host { if let ipv4Address = IPv4Address(ip) { return .ipv4(ipv4Address) } if let ipv6Address = IPv6Address(ip) { return .ipv6(ipv6Address) } preconditionFailure("invalid super server IP: \(ip)") } private func handleConnectionStateUpdate(_ state: NWConnection.State) { SDLLogger.log("[SDLSuperClient] new state: \(state)", category: .super) switch state { case .ready: Task { await self.readySignal.succeed(()) } case .failed(let error): let wrappedError = SDLSuperError.connectionFailed(error) Task { await self.readySignal.fail(wrappedError) } self.finishMessageStream(throwing: wrappedError) case .cancelled: let error = SDLSuperError.connectionCancelled Task { await self.readySignal.fail(error) } self.finishMessageStream(throwing: error) default: break } } private func readLoop() async throws { let frameParser = SDLSuperFrameParser(maxBufferSize: self.maxBufferSize) while true { try Task.checkCancellation() let data = try await Self.readOnce(connection: self.connection) let frames = try frameParser.parseFrames(data: data) for frame in frames { try Task.checkCancellation() if let message = SDLSuperCodec.decode(frame: frame) { self.messageContinuation.yield(message) } else { throw SDLSuperError.decodeError("invalid message") } } } } func send(type: SDLPacketType, data: Data) { guard connection.state == .ready else { return } var len = UInt16(data.count + 1).bigEndian var packet = Data(Data(bytes: &len, count: 2)) packet.append(type.rawValue) packet.append(data) connection.send(content: packet, completion: .contentProcessed { [weak self] error in if let error { SDLLogger.log("[SDLSuperClient] send data get error: \(error)", category: .super) self?.finishMessageStream(throwing: SDLSuperError.writeFailed(error)) self?.connection.cancel() } }) } private static func readOnce(connection: NWConnection) async throws -> Data { guard connection.state == .ready else { throw SDLSuperError.connectionCancelled } let readContinuation = OnceContinuation() return try await withTaskCancellationHandler { try await withCheckedThrowingContinuation { cont in readContinuation.set(cont) connection.receive(minimumIncompleteLength: 1, maximumLength: 64 * 1024) { data, _, isComplete, error in if let error { readContinuation.resume(throwing: error) return } if isComplete { readContinuation.resume(throwing: SDLSuperError.dataStreamClosed) } else { readContinuation.resume(returning: data ?? Data()) } } } } onCancel: { readContinuation.resume(throwing: CancellationError()) } } func stop() { guard self.markStopped() else { return } let connection = self.connection connection.stateUpdateHandler = nil connection.cancel() Task { await self.readySignal.fail(SDLSuperError.connectionCancelled) } self.finishMessageStream() SDLLogger.log("[SDLSuperClient] stopped", category: .super) } private func markStarted() -> Bool { self.stateLock.lock() defer { self.stateLock.unlock() } guard !self.isStarted, !self.isStopped else { return false } self.isStarted = true return true } private func markStopped() -> Bool { self.stateLock.lock() defer { self.stateLock.unlock() } guard !self.isStopped else { return false } self.isStopped = true return true } private func finishMessageStream(throwing error: Error? = nil) { self.stateLock.lock() let shouldFinish = !self.isMessageContinuationFinished if shouldFinish { self.isMessageContinuationFinished = true } self.stateLock.unlock() guard shouldFinish else { return } if let error { self.messageContinuation.finish(throwing: error) } else { self.messageContinuation.finish() } } deinit { SDLLogger.log("[SDLSuperClient] deinit", category: .super) } }