fix udp
This commit is contained in:
parent
f04196ef34
commit
a8a4cc507a
@ -170,7 +170,6 @@ actor SDLContextActor {
|
||||
let packetInboundActor = self.packetInboundActor
|
||||
let udpHoleService = SDLUDPHoleService(
|
||||
proberActor: self.proberActor,
|
||||
datagramSender: self.udpHoleServiceProxy.datagramSender,
|
||||
onEvent: udpHoleEventHandler,
|
||||
onData: { data in
|
||||
await packetInboundActor.handleData(data)
|
||||
|
||||
@ -263,10 +263,6 @@ actor PacketOutboundActor {
|
||||
}
|
||||
|
||||
private func sendPacket(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async {
|
||||
if type == .data, self.udpHoleServiceProxy.sendFast(type: type, data: data, remoteAddress: remoteAddress) {
|
||||
return
|
||||
}
|
||||
|
||||
await self.udpHoleServiceProxy.send(type: type, data: data, remoteAddress: remoteAddress)
|
||||
}
|
||||
}
|
||||
|
||||
@ -26,8 +26,8 @@ actor SDLUDPHole {
|
||||
private var state: State = .idle
|
||||
private let udpHoleHandler: SDLUDPHoleHandler
|
||||
|
||||
init(datagramSender: SDLUDPHoleDatagramSender) throws {
|
||||
self.udpHoleHandler = try SDLUDPHoleHandler(datagramSender: datagramSender)
|
||||
init() throws {
|
||||
self.udpHoleHandler = try SDLUDPHoleHandler()
|
||||
}
|
||||
|
||||
func start() async throws -> SocketAddress {
|
||||
@ -69,7 +69,6 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
||||
|
||||
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
||||
private var channel: Channel?
|
||||
private let datagramSender: SDLUDPHoleDatagramSender
|
||||
|
||||
private let locker = NSLock()
|
||||
|
||||
@ -78,8 +77,7 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
||||
private var isMessageContinuationFinished: Bool = false
|
||||
|
||||
// 启动函数
|
||||
init(datagramSender: SDLUDPHoleDatagramSender) throws {
|
||||
self.datagramSender = datagramSender
|
||||
init() throws {
|
||||
let (stream, continuation) = AsyncThrowingStream.makeStream(of: (SocketAddress, SDLHoleMessage).self, bufferingPolicy: .bufferingNewest(2048))
|
||||
self.messageStream = stream
|
||||
self.messageContinuation = continuation
|
||||
@ -99,7 +97,6 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
||||
}
|
||||
|
||||
self.channel = channel
|
||||
self.datagramSender.bind(channel, for: .v4)
|
||||
|
||||
return localAddress
|
||||
}
|
||||
@ -157,7 +154,6 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
||||
self.finishMessageContinuationIfNeed(throwing: nil)
|
||||
let channel = self.channel
|
||||
self.channel = nil
|
||||
self.datagramSender.bind(nil, for: .v4)
|
||||
try? channel?.close().wait()
|
||||
try? self.group.syncShutdownGracefully()
|
||||
|
||||
|
||||
@ -15,54 +15,6 @@ enum SDLUDPHoleKind: Equatable {
|
||||
}
|
||||
}
|
||||
|
||||
final class SDLUDPHoleDatagramSender: @unchecked Sendable {
|
||||
private let lock = NSLock()
|
||||
private var v4Channel: Channel?
|
||||
private var v6Channel: Channel?
|
||||
|
||||
func bind(_ channel: Channel?, for kind: SDLUDPHoleKind) {
|
||||
self.lock.lock()
|
||||
defer {
|
||||
self.lock.unlock()
|
||||
}
|
||||
|
||||
switch kind {
|
||||
case .v4:
|
||||
self.v4Channel = channel
|
||||
case .v6:
|
||||
self.v6Channel = channel
|
||||
}
|
||||
}
|
||||
|
||||
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) -> Bool {
|
||||
let channel: Channel?
|
||||
self.lock.lock()
|
||||
switch remoteAddress {
|
||||
case .v4:
|
||||
channel = self.v4Channel
|
||||
case .v6:
|
||||
channel = self.v6Channel
|
||||
default:
|
||||
channel = nil
|
||||
}
|
||||
self.lock.unlock()
|
||||
|
||||
guard let channel else {
|
||||
return false
|
||||
}
|
||||
|
||||
var buffer = channel.allocator.buffer(capacity: data.count + 1)
|
||||
buffer.writeBytes([type.rawValue])
|
||||
buffer.writeBytes(data)
|
||||
|
||||
let envelope = AddressedEnvelope<ByteBuffer>(remoteAddress: remoteAddress, data: buffer)
|
||||
channel.eventLoop.execute {
|
||||
channel.writeAndFlush(envelope, promise: nil)
|
||||
}
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
actor SDLUDPHoleService {
|
||||
enum Event {
|
||||
case ready(SocketAddress)
|
||||
@ -77,7 +29,6 @@ actor SDLUDPHoleService {
|
||||
private let proberActor: SDLNATProberActor
|
||||
private let onEvent: EventHandler
|
||||
private let onData: DataHandler
|
||||
private let datagramSender: SDLUDPHoleDatagramSender
|
||||
|
||||
private var udpHole: SDLUDPHole?
|
||||
private var udpHoleMonitorTask: Task<Void, Never>?
|
||||
@ -89,12 +40,10 @@ actor SDLUDPHoleService {
|
||||
|
||||
init(
|
||||
proberActor: SDLNATProberActor,
|
||||
datagramSender: SDLUDPHoleDatagramSender,
|
||||
onEvent: @escaping EventHandler,
|
||||
onData: @escaping DataHandler
|
||||
) {
|
||||
self.proberActor = proberActor
|
||||
self.datagramSender = datagramSender
|
||||
self.onEvent = onEvent
|
||||
self.onData = onData
|
||||
}
|
||||
@ -173,7 +122,7 @@ actor SDLUDPHoleService {
|
||||
}
|
||||
|
||||
private func runV4() async throws {
|
||||
let udpHole = try SDLUDPHole(datagramSender: self.datagramSender)
|
||||
let udpHole = try SDLUDPHole()
|
||||
let localAddress = try await udpHole.start()
|
||||
self.udpHole = udpHole
|
||||
self.localAddress = localAddress
|
||||
@ -250,7 +199,7 @@ actor SDLUDPHoleService {
|
||||
}
|
||||
|
||||
private func runV6() async throws {
|
||||
let udpHoleV6 = try SDLUDPHoleV6(datagramSender: self.datagramSender)
|
||||
let udpHoleV6 = try SDLUDPHoleV6()
|
||||
let localAddress = try udpHoleV6.start()
|
||||
self.udpHoleV6 = udpHoleV6
|
||||
|
||||
@ -306,7 +255,6 @@ actor SDLUDPHoleService {
|
||||
actor SDLUDPHoleServiceProxy {
|
||||
typealias ControlEventHandler = @Sendable (SDLUDPHoleService.Event) async -> Void
|
||||
|
||||
nonisolated let datagramSender = SDLUDPHoleDatagramSender()
|
||||
private var udpHoleService: SDLUDPHoleService?
|
||||
private var generation: UInt64 = 0
|
||||
|
||||
@ -341,10 +289,6 @@ actor SDLUDPHoleServiceProxy {
|
||||
await self.udpHoleService?.send(type: type, data: data, remoteAddress: remoteAddress)
|
||||
}
|
||||
|
||||
nonisolated func sendFast(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) -> Bool {
|
||||
self.datagramSender.send(type: type, data: data, remoteAddress: remoteAddress)
|
||||
}
|
||||
|
||||
private func handleEvent(_ event: SDLUDPHoleService.Event, generation: UInt64, onControlEvent: ControlEventHandler) async {
|
||||
guard generation == self.generation else {
|
||||
return
|
||||
|
||||
@ -25,7 +25,6 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
||||
|
||||
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
||||
private var channel: Channel?
|
||||
private let datagramSender: SDLUDPHoleDatagramSender
|
||||
|
||||
public let messageStream: AsyncStream<(SocketAddress, SDLHoleMessage)>
|
||||
private let messageContinuation: AsyncStream<(SocketAddress, SDLHoleMessage)>.Continuation
|
||||
@ -35,8 +34,7 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
||||
private let eventContinuation: AsyncStream<HoleEvent>.Continuation
|
||||
|
||||
// 启动函数
|
||||
init(datagramSender: SDLUDPHoleDatagramSender) throws {
|
||||
self.datagramSender = datagramSender
|
||||
init() throws {
|
||||
let (stream, continuation) = AsyncStream.makeStream(of: (SocketAddress, SDLHoleMessage).self, bufferingPolicy: .bufferingNewest(2048))
|
||||
self.messageStream = stream
|
||||
self.messageContinuation = continuation
|
||||
@ -56,7 +54,6 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
||||
// 绑定到IPv6通配地址,只处理IPv6流量
|
||||
let channel = try bootstrap.bind(host: "::", port: 0).wait()
|
||||
self.channel = channel
|
||||
self.datagramSender.bind(channel, for: .v6)
|
||||
|
||||
return channel.localAddress
|
||||
}
|
||||
@ -124,7 +121,6 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
||||
|
||||
let channel = self.channel
|
||||
self.channel = nil
|
||||
self.datagramSender.bind(nil, for: .v6)
|
||||
try? channel?.close().wait()
|
||||
|
||||
try? self.group.syncShutdownGracefully()
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user