fix dnsClient

This commit is contained in:
anlicheng 2026-05-04 21:12:22 +08:00
parent ad6f7a8c42
commit 288c58d1a7
2 changed files with 135 additions and 152 deletions

View File

@ -56,7 +56,6 @@ actor SDLContextActor {
// dnsclient
private var dnsClient: DNSCloudClient?
private var dnsWorker: Task<Void, Never>?
// Localdnsclient
private var dnsLocalClient: DNSLocalClient?
@ -124,10 +123,20 @@ actor SDLContextActor {
// arp
await self.puncherActor.start()
await self.arpServer.start()
await self.startDnsClient()
await self.startDnsLocalClient()
// udp
await self.supervisor.addWorker(name: "dnsClient") {
SDLLogger.log("[SDLContext] dnsClient running!!!!")
try await self.startDnsClient()
SDLLogger.log("[SDLContext] dnsClient closed!!!!")
}
await self.supervisor.addWorker(name: "dnsLocalClient") {
SDLLogger.log("[SDLContext] dnsLocalClient running!!!!")
try await self.startDnsLocalClient()
SDLLogger.log("[SDLContext] dnsLocalClient closed!!!!")
}
// udp
await self.supervisor.addWorker(name: "udpHole") {
SDLLogger.log("[SDLContext] udp running!!!!")
try await self.startUDPHole()
@ -294,26 +303,52 @@ actor SDLContextActor {
SDLLogger.log("[SDLContext] tunnelAppNotifier ready")
}
private func startDnsClient() async {
self.dnsWorker?.cancel()
self.dnsWorker = nil
private func startDnsClient() async throws {
// dns
let dnsClient = DNSCloudClient(host: self.config.serverHost, port: 15353)
await dnsClient.start()
dnsClient.start()
SDLLogger.log("[SDLContext] dnsClient started")
self.dnsClient = dnsClient
let packetFlow = dnsClient.packetFlow
self.dnsWorker = Task.detached {
//
for await packet in packetFlow {
if Task.isCancelled {
break
defer {
self.dnsClient?.stop()
self.dnsClient = nil
}
try await withThrowingTaskGroup { group in
defer {
group.cancelAll()
}
group.addTask {
for await packet in dnsClient.packetFlow {
try Task.checkCancellation()
let nePacket = NEPacket(data: packet, protocolFamily: 2)
self.provider.packetFlow.writePacketObjects([nePacket])
}
}
group.addTask {
for await event in dnsClient.eventStream {
try Task.checkCancellation()
switch event {
case .failed(let error):
SDLLogger.log("[SDLContext] dnsClient failed with error: \(error)")
throw error
case .cancelled:
SDLLogger.log("[SDLContext] dnsClient cancelled")
return
case .sendFailed(let error):
SDLLogger.log("[SDLContext] dnsClient sendFailed with error: \(error)")
throw error
}
}
}
try await group.next()
}
}
private func startDnsLocalClient() async {

View File

@ -7,7 +7,14 @@
import Foundation
import Network
actor DNSCloudClient {
final class DNSCloudClient {
enum Event {
case failed(Error)
case cancelled
case sendFailed(Error)
}
private enum State {
case idle
case running
@ -15,6 +22,7 @@ actor DNSCloudClient {
}
private var state: State = .idle
private var connection: NWConnection?
private var receiveTask: Task<Void, Never>?
private let dnsServerAddress: NWEndpoint
@ -22,37 +30,28 @@ actor DNSCloudClient {
// DNS
public let packetFlow: AsyncStream<Data>
private let packetContinuation: AsyncStream<Data>.Continuation
private var didFinishPacketFlow = false
//
private let closeStream: AsyncStream<Void>
private let closeContinuation: AsyncStream<Void>.Continuation
private var didFinishCloseStream = false
// Connection
public let eventStream: AsyncStream<Event>
private let eventContinuation: AsyncStream<Event>.Continuation
/// - Parameter host: sn-server ( "8.8.8.8")
/// - Parameter port: ( 53)
init(host: String, port: UInt16 ) {
self.dnsServerAddress = .hostPort(host: NWEndpoint.Host(host), port: NWEndpoint.Port(integerLiteral: port))
let (packetStream, packetContinuation) = AsyncStream.makeStream(of: Data.self, bufferingPolicy: .bufferingNewest(256))
self.packetFlow = packetStream
self.packetContinuation = packetContinuation
let packetPair = AsyncStream.makeStream(of: Data.self, bufferingPolicy: .bufferingNewest(256))
self.packetFlow = packetPair.stream
self.packetContinuation = packetPair.continuation
let (closeStream, closeContinuation) = AsyncStream.makeStream(of: Void.self, bufferingPolicy: .bufferingNewest(1))
self.closeStream = closeStream
self.closeContinuation = closeContinuation
let eventPair = AsyncStream.makeStream(of: Event.self)
self.eventStream = eventPair.stream
self.eventContinuation = eventPair.continuation
}
func start() {
guard case .idle = self.state else {
return
}
self.state = .running
// 1.
let parameters = NWParameters.udp
// TUN NE TUN .other
parameters.prohibitedInterfaceTypes = [.other]
// 2. pathSelectionOptions
@ -60,20 +59,74 @@ actor DNSCloudClient {
// 2.
let connection = NWConnection(to: self.dnsServerAddress, using: parameters)
self.connection = connection
connection.stateUpdateHandler = { [weak self] state in
Task {
await self?.handleConnectionStateUpdate(state, for: connection)
self?.handleConnectionStateUpdate(state, for: connection)
}
}
//
connection.start(queue: .global())
self.connection = connection
}
public func waitClose() async {
for await _ in self.closeStream { }
/// DNS TUN IP
func forward(ipPacketData: Data) {
guard let connection = self.connection, connection.state == .ready else {
return
}
connection.send(content: ipPacketData, completion: .contentProcessed { error in
if let error = error {
self.eventContinuation.yield(.sendFailed(error))
}
})
}
func stop() {
guard self.state != .stopped else {
return
}
self.state = .stopped
self.receiveTask?.cancel()
self.receiveTask = nil
self.connection?.cancel()
self.connection = nil
self.packetContinuation.finish()
self.eventContinuation.finish()
}
private func handleConnectionStateUpdate(_ state: NWConnection.State, for connection: NWConnection) {
switch state {
case .ready:
SDLLogger.log("[DNSClient] Connection ready", for: .debug)
self.startReceiveTask(for: connection)
self.state = .running
case .failed(let error):
self.eventContinuation.yield(.failed(error))
case .cancelled:
self.eventContinuation.yield(.cancelled)
default:
break
}
}
private func startReceiveTask(for connection: NWConnection) {
guard self.receiveTask == nil else {
return
}
let stream = Self.makeReceiveStream(for: connection)
self.receiveTask = Task { [weak self] in
for await data in stream {
if Task.isCancelled {
break
}
self?.packetContinuation.yield(data)
}
}
}
///
@ -98,109 +151,4 @@ actor DNSCloudClient {
}
}
/// DNS TUN IP
func forward(ipPacketData: Data) {
guard case .running = self.state, let connection = self.connection, connection.state == .ready else {
return
}
connection.send(content: ipPacketData, completion: .contentProcessed { error in
if let error = error {
SDLLogger.log("[DNSClient] Send error: \(error)", for: .debug)
}
})
}
func stop() {
guard self.state != .stopped else {
return
}
self.state = .stopped
self.receiveTask?.cancel()
self.receiveTask = nil
self.connection?.cancel()
self.connection = nil
self.finishPacketFlowIfNeeded()
self.finishCloseStreamIfNeeded()
}
private func handleConnectionStateUpdate(_ state: NWConnection.State, for connection: NWConnection) {
guard case .running = self.state else {
return
}
switch state {
case .ready:
SDLLogger.log("[DNSClient] Connection ready", for: .debug)
self.startReceiveTask(for: connection)
case .failed(let error):
SDLLogger.log("[DNSClient] Connection failed: \(error)", for: .debug)
self.stop()
case .cancelled:
self.stop()
default:
break
}
}
private func startReceiveTask(for connection: NWConnection) {
guard self.receiveTask == nil else {
return
}
let stream = Self.makeReceiveStream(for: connection)
self.receiveTask = Task { [weak self] in
for await data in stream {
guard let self else {
break
}
await self.handleReceivedPacket(data)
}
await self?.didFinishReceiving(for: connection)
}
}
private func handleReceivedPacket(_ data: Data) {
guard case .running = self.state else {
return
}
self.packetContinuation.yield(data)
}
private func didFinishReceiving(for connection: NWConnection) {
guard case .running = self.state else {
return
}
if self.connection === connection, connection.state != .ready {
self.stop()
} else {
self.receiveTask = nil
}
}
private func finishPacketFlowIfNeeded() {
guard !self.didFinishPacketFlow else {
return
}
self.didFinishPacketFlow = true
self.packetContinuation.finish()
}
private func finishCloseStreamIfNeeded() {
guard !self.didFinishCloseStream else {
return
}
self.didFinishCloseStream = true
self.closeContinuation.finish()
}
deinit {
self.connection?.cancel()
}
}