fix udpHole
This commit is contained in:
parent
783c1329e1
commit
55701085e3
@ -53,8 +53,6 @@ actor SDLContextActor {
|
|||||||
private var udpHoleLocalAddress: SocketAddress?
|
private var udpHoleLocalAddress: SocketAddress?
|
||||||
|
|
||||||
private var udpHoleV6: SDLUDPHoleV6?
|
private var udpHoleV6: SDLUDPHoleV6?
|
||||||
private var udpHoleV6Workers: [Task<Void, Never>]?
|
|
||||||
private var udpHoleV6LocalAddress: SocketAddress?
|
|
||||||
|
|
||||||
// dns的client对象
|
// dns的client对象
|
||||||
private var dnsClient: DNSCloudClient?
|
private var dnsClient: DNSCloudClient?
|
||||||
@ -149,9 +147,8 @@ actor SDLContextActor {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// await self.supervisor.addWorker(name: "udpHoleV6") {
|
// await self.supervisor.addWorker(name: "udpHoleV6") {
|
||||||
// let udpHoleV6 = try await self.startUDPHoleV6()
|
|
||||||
// SDLLogger.log("[SDLContext] udp v6 running!!!!")
|
// SDLLogger.log("[SDLContext] udp v6 running!!!!")
|
||||||
// try await udpHoleV6.waitClose()
|
// try await self.startUDPHoleV6()
|
||||||
// SDLLogger.log("[SDLContext] udp v6 closed!!!!")
|
// SDLLogger.log("[SDLContext] udp v6 closed!!!!")
|
||||||
// }
|
// }
|
||||||
}
|
}
|
||||||
@ -357,6 +354,7 @@ actor SDLContextActor {
|
|||||||
self.udpHoleLocalAddress = localAddress
|
self.udpHoleLocalAddress = localAddress
|
||||||
|
|
||||||
defer {
|
defer {
|
||||||
|
self.udpHole?.stop()
|
||||||
self.udpHole = nil
|
self.udpHole = nil
|
||||||
self.udpHoleLocalAddress = nil
|
self.udpHoleLocalAddress = nil
|
||||||
}
|
}
|
||||||
@ -371,8 +369,17 @@ actor SDLContextActor {
|
|||||||
try Task.checkCancellation()
|
try Task.checkCancellation()
|
||||||
|
|
||||||
switch message.inboundMessage {
|
switch message.inboundMessage {
|
||||||
case .control(let controlMessage):
|
case .control(let message):
|
||||||
await self.handleHoleControlMessage(controlMessage, localAddress: localAddress, remoteAddress: remoteAddress, source: .v4)
|
switch message {
|
||||||
|
case .stunReply(_):
|
||||||
|
SDLLogger.log("[SDLContext] get a stunReply", for: .debug)
|
||||||
|
case .stunProbeReply(let probeReply):
|
||||||
|
await self.proberActor.handleProbeReply(localAddress: localAddress, reply: probeReply)
|
||||||
|
case .register(let register):
|
||||||
|
try? await self.handleRegister(remoteAddress: remoteAddress, register: register, source: .v4)
|
||||||
|
case .registerAck(let registerAck):
|
||||||
|
await self.handleRegisterAck(remoteAddress: remoteAddress, registerAck: registerAck, source: .v6)
|
||||||
|
}
|
||||||
case .data(let data):
|
case .data(let data):
|
||||||
try? await self.handleHoleData(data: data)
|
try? await self.handleHoleData(data: data)
|
||||||
}
|
}
|
||||||
@ -399,26 +406,64 @@ actor SDLContextActor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private func startUDPHoleV6() async throws -> SDLUDPHoleV6 {
|
private func startUDPHoleV6() async throws {
|
||||||
self.udpHoleV6Workers?.forEach {$0.cancel()}
|
|
||||||
self.udpHoleV6Workers = nil
|
|
||||||
|
|
||||||
// 启动udp服务器
|
// 启动udp服务器
|
||||||
let udpHoleV6 = try SDLUDPHoleV6()
|
let udpHoleV6 = try SDLUDPHoleV6()
|
||||||
let localAddress = try udpHoleV6.start()
|
let localAddress = try udpHoleV6.start()
|
||||||
SDLLogger.log("[SDLContext] udpHoleV6 started, on address: \(localAddress)")
|
self.udpHoleV6 = udpHoleV6
|
||||||
|
|
||||||
// 处理消息流
|
if let localAddress {
|
||||||
let messageStream = udpHoleV6.messageStream
|
SDLLogger.log("[SDLContext] udpHoleV6 started, on address: \(localAddress)")
|
||||||
let messageTask = Task.detached {
|
} else {
|
||||||
await self.consumeUDPHoleMessages(stream: messageStream, localAddress: localAddress, source: .v6)
|
SDLLogger.log("[SDLContext] udpHoleV6 started, no local address")
|
||||||
}
|
}
|
||||||
|
|
||||||
self.udpHoleV6 = udpHoleV6
|
defer {
|
||||||
self.udpHoleV6LocalAddress = localAddress
|
self.udpHoleV6?.stop()
|
||||||
self.udpHoleV6Workers = [messageTask]
|
self.udpHoleV6 = nil
|
||||||
|
}
|
||||||
|
|
||||||
return udpHoleV6
|
try await withThrowingTaskGroup { group in
|
||||||
|
defer {
|
||||||
|
group.cancelAll()
|
||||||
|
}
|
||||||
|
|
||||||
|
// 处理消息流
|
||||||
|
group.addTask {
|
||||||
|
for await (remoteAddress, message) in udpHoleV6.messageStream {
|
||||||
|
try Task.checkCancellation()
|
||||||
|
|
||||||
|
switch message.inboundMessage {
|
||||||
|
case .control(let message):
|
||||||
|
switch message {
|
||||||
|
case .register(let register):
|
||||||
|
try? await self.handleRegister(remoteAddress: remoteAddress, register: register, source: .v6)
|
||||||
|
case .registerAck(let registerAck):
|
||||||
|
await self.handleRegisterAck(remoteAddress: remoteAddress, registerAck: registerAck, source: .v6)
|
||||||
|
default:
|
||||||
|
()
|
||||||
|
}
|
||||||
|
case .data(let data):
|
||||||
|
try? await self.handleHoleData(data: data)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
group.addTask {
|
||||||
|
for await event in udpHoleV6.eventStream {
|
||||||
|
try Task.checkCancellation()
|
||||||
|
|
||||||
|
switch event {
|
||||||
|
case .ready:
|
||||||
|
SDLLogger.log("[SDLContext] udpHoleV6 ready")
|
||||||
|
case .closed, .errorCaught:
|
||||||
|
throw SDLContextError.udpHoleClosed
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
try await group.next()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 处理context的停止问题
|
// 处理context的停止问题
|
||||||
@ -435,14 +480,12 @@ actor SDLContextActor {
|
|||||||
|
|
||||||
self.flowSessionManager.clear()
|
self.flowSessionManager.clear()
|
||||||
|
|
||||||
|
self.udpHole?.stop()
|
||||||
self.udpHole = nil
|
self.udpHole = nil
|
||||||
self.udpHoleLocalAddress = nil
|
self.udpHoleLocalAddress = nil
|
||||||
|
|
||||||
self.udpHoleV6Workers?.forEach { $0.cancel() }
|
|
||||||
self.udpHoleV6Workers = nil
|
|
||||||
self.udpHoleV6?.stop()
|
self.udpHoleV6?.stop()
|
||||||
self.udpHoleV6 = nil
|
self.udpHoleV6 = nil
|
||||||
self.udpHoleV6LocalAddress = nil
|
|
||||||
|
|
||||||
self.quicClient?.stop()
|
self.quicClient?.stop()
|
||||||
self.quicClient = nil
|
self.quicClient = nil
|
||||||
@ -682,7 +725,6 @@ actor SDLContextActor {
|
|||||||
self.udpHole = nil
|
self.udpHole = nil
|
||||||
self.udpHoleLocalAddress = nil
|
self.udpHoleLocalAddress = nil
|
||||||
self.udpHoleV6 = nil
|
self.udpHoleV6 = nil
|
||||||
self.udpHoleV6LocalAddress = nil
|
|
||||||
self.dnsClient = nil
|
self.dnsClient = nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -800,40 +842,6 @@ extension SDLContextActor {
|
|||||||
// 处理从Hole收到的数据
|
// 处理从Hole收到的数据
|
||||||
extension SDLContextActor {
|
extension SDLContextActor {
|
||||||
|
|
||||||
private func consumeUDPHoleMessages(stream: AsyncStream<(SocketAddress, SDLHoleMessage)>, localAddress: SocketAddress, source: UDPHoleKind) async {
|
|
||||||
for await (remoteAddress, message) in stream {
|
|
||||||
if Task.isCancelled {
|
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
switch message.inboundMessage {
|
|
||||||
case .control(let controlMessage):
|
|
||||||
await self.handleHoleControlMessage(controlMessage, localAddress: localAddress, remoteAddress: remoteAddress, source: source)
|
|
||||||
case .data(let data):
|
|
||||||
try? await self.handleHoleData(data: data)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private func handleHoleControlMessage(_ message: SDLHoleControlMessage, localAddress: SocketAddress, remoteAddress: SocketAddress, source: UDPHoleKind) async {
|
|
||||||
switch message {
|
|
||||||
case .stunReply(_):
|
|
||||||
guard source == .v4 else {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
SDLLogger.log("[SDLContext] get a stunReply", for: .debug)
|
|
||||||
case .stunProbeReply(let probeReply):
|
|
||||||
guard source == .v4 else {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
await self.proberActor.handleProbeReply(localAddress: localAddress, reply: probeReply)
|
|
||||||
case .register(let register):
|
|
||||||
try? await self.handleRegister(remoteAddress: remoteAddress, register: register, source: source)
|
|
||||||
case .registerAck(let registerAck):
|
|
||||||
await self.handleRegisterAck(remoteAddress: remoteAddress, registerAck: registerAck, source: source)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private func handleRegister(remoteAddress: SocketAddress, register: SDLRegister, source: UDPHoleKind) async throws {
|
private func handleRegister(remoteAddress: SocketAddress, register: SDLRegister, source: UDPHoleKind) async throws {
|
||||||
let networkAddr = config.networkAddress
|
let networkAddr = config.networkAddress
|
||||||
SDLLogger.log("[SDLContext] register packet: \(register), network_address: \(networkAddr)")
|
SDLLogger.log("[SDLContext] register packet: \(register), network_address: \(networkAddr)")
|
||||||
|
|||||||
@ -24,6 +24,8 @@ final class SDLUDPHole: ChannelInboundHandler {
|
|||||||
case errorCaught
|
case errorCaught
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private var isStopped: Bool = false
|
||||||
|
|
||||||
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
||||||
private var channel: Channel?
|
private var channel: Channel?
|
||||||
|
|
||||||
@ -113,8 +115,14 @@ final class SDLUDPHole: ChannelInboundHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
deinit {
|
func stop() {
|
||||||
SDLLogger.log("[SDLUDPHole] deinit", for: .debug)
|
guard !self.isStopped else {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
self.isStopped = true
|
||||||
|
|
||||||
|
SDLLogger.log("[SDLUDPHole] stop", for: .debug)
|
||||||
self.messageContinuation.finish()
|
self.messageContinuation.finish()
|
||||||
self.eventContinuation.finish()
|
self.eventContinuation.finish()
|
||||||
self.channel = nil
|
self.channel = nil
|
||||||
|
|||||||
@ -14,30 +14,37 @@ import SwiftProtobuf
|
|||||||
final class SDLUDPHoleV6: ChannelInboundHandler {
|
final class SDLUDPHoleV6: ChannelInboundHandler {
|
||||||
typealias InboundIn = AddressedEnvelope<ByteBuffer>
|
typealias InboundIn = AddressedEnvelope<ByteBuffer>
|
||||||
|
|
||||||
private enum State: Equatable {
|
// 事件
|
||||||
case idle
|
enum HoleEvent {
|
||||||
case ready
|
case ready
|
||||||
case stopping
|
case closed
|
||||||
case stopped
|
case errorCaught
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private var isStopped: Bool = false
|
||||||
|
|
||||||
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
||||||
private var channel: Channel?
|
private var channel: Channel?
|
||||||
private var closeFuture: EventLoopFuture<Void>?
|
|
||||||
private var state: State = .idle
|
|
||||||
private var didFinishMessageStream: Bool = false
|
|
||||||
|
|
||||||
public let messageStream: AsyncStream<(SocketAddress, SDLHoleMessage)>
|
public let messageStream: AsyncStream<(SocketAddress, SDLHoleMessage)>
|
||||||
private let messageContinuation: AsyncStream<(SocketAddress, SDLHoleMessage)>.Continuation
|
private let messageContinuation: AsyncStream<(SocketAddress, SDLHoleMessage)>.Continuation
|
||||||
|
|
||||||
|
// 事件相关逻辑
|
||||||
|
public let eventStream: AsyncStream<HoleEvent>
|
||||||
|
private let eventContinuation: AsyncStream<HoleEvent>.Continuation
|
||||||
|
|
||||||
// 启动函数
|
// 启动函数
|
||||||
init() throws {
|
init() throws {
|
||||||
let (stream, continuation) = AsyncStream.makeStream(of: (SocketAddress, SDLHoleMessage).self, bufferingPolicy: .bufferingNewest(2048))
|
let (stream, continuation) = AsyncStream.makeStream(of: (SocketAddress, SDLHoleMessage).self, bufferingPolicy: .bufferingNewest(2048))
|
||||||
self.messageStream = stream
|
self.messageStream = stream
|
||||||
self.messageContinuation = continuation
|
self.messageContinuation = continuation
|
||||||
|
|
||||||
|
let eventPair = AsyncStream.makeStream(of: HoleEvent.self)
|
||||||
|
self.eventStream = eventPair.stream
|
||||||
|
self.eventContinuation = eventPair.continuation
|
||||||
}
|
}
|
||||||
|
|
||||||
func start() throws -> SocketAddress {
|
func start() throws -> SocketAddress? {
|
||||||
let bootstrap = DatagramBootstrap(group: group)
|
let bootstrap = DatagramBootstrap(group: group)
|
||||||
.channelOption(ChannelOptions.socketOption(.so_reuseaddr), value: 1)
|
.channelOption(ChannelOptions.socketOption(.so_reuseaddr), value: 1)
|
||||||
.channelInitializer { channel in
|
.channelInitializer { channel in
|
||||||
@ -47,51 +54,13 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
|||||||
// 绑定到IPv6通配地址,只处理IPv6流量
|
// 绑定到IPv6通配地址,只处理IPv6流量
|
||||||
let channel = try bootstrap.bind(host: "::", port: 0).wait()
|
let channel = try bootstrap.bind(host: "::", port: 0).wait()
|
||||||
self.channel = channel
|
self.channel = channel
|
||||||
self.closeFuture = channel.closeFuture
|
|
||||||
self.state = .ready
|
|
||||||
precondition(channel.localAddress != nil, "UDP v6 channel has no localAddress after bind")
|
|
||||||
|
|
||||||
return channel.localAddress!
|
return channel.localAddress
|
||||||
}
|
|
||||||
|
|
||||||
func waitClose() async throws {
|
|
||||||
switch self.state {
|
|
||||||
case .idle:
|
|
||||||
SDLLogger.log("[SDLUDPHoleV6] waitClose11", for: .debug)
|
|
||||||
return
|
|
||||||
case .ready, .stopping, .stopped:
|
|
||||||
guard let closeFuture = self.closeFuture else {
|
|
||||||
SDLLogger.log("[SDLUDPHoleV6] waitClose22", for: .debug)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
try await closeFuture.get()
|
|
||||||
SDLLogger.log("[SDLUDPHoleV6] waitClose33", for: .debug)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func stop() {
|
|
||||||
switch self.state {
|
|
||||||
case .stopping, .stopped:
|
|
||||||
return
|
|
||||||
case .idle:
|
|
||||||
self.state = .stopped
|
|
||||||
self.finishMessageStream()
|
|
||||||
return
|
|
||||||
case .ready:
|
|
||||||
self.state = .stopping
|
|
||||||
}
|
|
||||||
|
|
||||||
self.finishMessageStream()
|
|
||||||
self.channel?.close(promise: nil)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// --MARK: ChannelInboundHandler delegate
|
// --MARK: ChannelInboundHandler delegate
|
||||||
|
|
||||||
func channelRead(context: ChannelHandlerContext, data: NIOAny) {
|
func channelRead(context: ChannelHandlerContext, data: NIOAny) {
|
||||||
guard case .ready = self.state else {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
let envelope = unwrapInboundIn(data)
|
let envelope = unwrapInboundIn(data)
|
||||||
|
|
||||||
var buffer = envelope.data
|
var buffer = envelope.data
|
||||||
@ -113,23 +82,19 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func channelInactive(context: ChannelHandlerContext) {
|
func channelInactive(context: ChannelHandlerContext) {
|
||||||
self.finishMessageStream()
|
SDLLogger.log("[SDLUDPHoleV6] channelInactive", for: .debug)
|
||||||
self.channel = nil
|
self.eventContinuation.yield(.closed)
|
||||||
self.state = .stopped
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func errorCaught(context: ChannelHandlerContext, error: any Error) {
|
func errorCaught(context: ChannelHandlerContext, error: any Error) {
|
||||||
SDLLogger.log("[SDLUDPHoleV6] channel error: \(error)", for: .debug)
|
SDLLogger.log("[SDLUDPHoleV6] channel error: \(error)", for: .debug)
|
||||||
self.finishMessageStream()
|
|
||||||
if self.state != .stopped {
|
|
||||||
self.state = .stopping
|
|
||||||
}
|
|
||||||
context.close(promise: nil)
|
context.close(promise: nil)
|
||||||
|
self.eventContinuation.yield(.errorCaught)
|
||||||
}
|
}
|
||||||
|
|
||||||
// MARK: 处理写入逻辑
|
// MARK: 处理写入逻辑
|
||||||
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) {
|
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) {
|
||||||
guard case .ready = self.state, let channel = self.channel else {
|
guard let channel = self.channel else {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -143,17 +108,18 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private func finishMessageStream() {
|
func stop() {
|
||||||
guard !self.didFinishMessageStream else {
|
guard !self.isStopped else {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
self.didFinishMessageStream = true
|
self.isStopped = true
|
||||||
self.messageContinuation.finish()
|
SDLLogger.log("[SDLUDPHoleV6] stop", for: .debug)
|
||||||
}
|
|
||||||
|
self.messageContinuation.finish()
|
||||||
|
self.eventContinuation.finish()
|
||||||
|
self.channel = nil
|
||||||
|
|
||||||
deinit {
|
|
||||||
self.stop()
|
|
||||||
try? self.group.syncShutdownGracefully()
|
try? self.group.syncShutdownGracefully()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user