fix udp
This commit is contained in:
parent
1042b5c724
commit
f04196ef34
@ -164,12 +164,18 @@ actor SDLContextActor {
|
|||||||
await self.packetOutboundActor.updateDNSService(dnsService)
|
await self.packetOutboundActor.updateDNSService(dnsService)
|
||||||
await dnsService.start()
|
await dnsService.start()
|
||||||
|
|
||||||
await self.udpHoleServiceProxy.bindInbound(self.packetInboundActor) { [weak self] event in
|
let udpHoleEventHandler = await self.udpHoleServiceProxy.makeEventHandler { [weak self] event in
|
||||||
await self?.handleUDPHoleControlEvent(event)
|
await self?.handleUDPHoleControlEvent(event)
|
||||||
}
|
}
|
||||||
|
let packetInboundActor = self.packetInboundActor
|
||||||
let udpHoleEventHandler = await self.udpHoleServiceProxy.makeEventHandler()
|
let udpHoleService = SDLUDPHoleService(
|
||||||
let udpHoleService = SDLUDPHoleService(proberActor: self.proberActor, onEvent: udpHoleEventHandler)
|
proberActor: self.proberActor,
|
||||||
|
datagramSender: self.udpHoleServiceProxy.datagramSender,
|
||||||
|
onEvent: udpHoleEventHandler,
|
||||||
|
onData: { data in
|
||||||
|
await packetInboundActor.handleData(data)
|
||||||
|
}
|
||||||
|
)
|
||||||
await self.udpHoleServiceProxy.replace(udpHoleService)
|
await self.udpHoleServiceProxy.replace(udpHoleService)
|
||||||
await udpHoleService.start(includeV6: false)
|
await udpHoleService.start(includeV6: false)
|
||||||
|
|
||||||
@ -469,9 +475,7 @@ extension SDLContextActor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private func handleUDPHolePacket(remoteAddress: SocketAddress, message: SDLHoleMessage, source: SDLUDPHoleKind) async {
|
private func handleUDPHolePacket(remoteAddress: SocketAddress, message: SDLHoleControlMessage, source: SDLUDPHoleKind) async {
|
||||||
switch message {
|
|
||||||
case .control(let message):
|
|
||||||
switch message {
|
switch message {
|
||||||
case .stunReply(_), .stunProbeReply(_):
|
case .stunReply(_), .stunProbeReply(_):
|
||||||
SDLLogger.log("[SDLContext] get a stun reply", for: .debug)
|
SDLLogger.log("[SDLContext] get a stun reply", for: .debug)
|
||||||
@ -480,9 +484,6 @@ extension SDLContextActor {
|
|||||||
case .registerAck(let registerAck):
|
case .registerAck(let registerAck):
|
||||||
await self.handleRegisterAck(remoteAddress: remoteAddress, registerAck: registerAck, source: source)
|
await self.handleRegisterAck(remoteAddress: remoteAddress, registerAck: registerAck, source: source)
|
||||||
}
|
}
|
||||||
case .data:
|
|
||||||
SDLLogger.log("[SDLContext] unexpected data packet in control path", for: .debug)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private func handleRegister(remoteAddress: SocketAddress, register: SDLRegister, source: SDLUDPHoleKind) async throws {
|
private func handleRegister(remoteAddress: SocketAddress, register: SDLRegister, source: SDLUDPHoleKind) async throws {
|
||||||
|
|||||||
@ -263,6 +263,10 @@ actor PacketOutboundActor {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private func sendPacket(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async {
|
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)
|
await self.udpHoleServiceProxy.send(type: type, data: data, remoteAddress: remoteAddress)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -26,8 +26,8 @@ actor SDLUDPHole {
|
|||||||
private var state: State = .idle
|
private var state: State = .idle
|
||||||
private let udpHoleHandler: SDLUDPHoleHandler
|
private let udpHoleHandler: SDLUDPHoleHandler
|
||||||
|
|
||||||
init() throws {
|
init(datagramSender: SDLUDPHoleDatagramSender) throws {
|
||||||
self.udpHoleHandler = try SDLUDPHoleHandler()
|
self.udpHoleHandler = try SDLUDPHoleHandler(datagramSender: datagramSender)
|
||||||
}
|
}
|
||||||
|
|
||||||
func start() async throws -> SocketAddress {
|
func start() async throws -> SocketAddress {
|
||||||
@ -69,6 +69,7 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
|||||||
|
|
||||||
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
||||||
private var channel: Channel?
|
private var channel: Channel?
|
||||||
|
private let datagramSender: SDLUDPHoleDatagramSender
|
||||||
|
|
||||||
private let locker = NSLock()
|
private let locker = NSLock()
|
||||||
|
|
||||||
@ -77,7 +78,8 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
|||||||
private var isMessageContinuationFinished: Bool = false
|
private var isMessageContinuationFinished: Bool = false
|
||||||
|
|
||||||
// 启动函数
|
// 启动函数
|
||||||
init() throws {
|
init(datagramSender: SDLUDPHoleDatagramSender) throws {
|
||||||
|
self.datagramSender = datagramSender
|
||||||
let (stream, continuation) = AsyncThrowingStream.makeStream(of: (SocketAddress, SDLHoleMessage).self, bufferingPolicy: .bufferingNewest(2048))
|
let (stream, continuation) = AsyncThrowingStream.makeStream(of: (SocketAddress, SDLHoleMessage).self, bufferingPolicy: .bufferingNewest(2048))
|
||||||
self.messageStream = stream
|
self.messageStream = stream
|
||||||
self.messageContinuation = continuation
|
self.messageContinuation = continuation
|
||||||
@ -97,6 +99,7 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
|||||||
}
|
}
|
||||||
|
|
||||||
self.channel = channel
|
self.channel = channel
|
||||||
|
self.datagramSender.bind(channel, for: .v4)
|
||||||
|
|
||||||
return localAddress
|
return localAddress
|
||||||
}
|
}
|
||||||
@ -154,6 +157,7 @@ private final class SDLUDPHoleHandler: ChannelInboundHandler {
|
|||||||
self.finishMessageContinuationIfNeed(throwing: nil)
|
self.finishMessageContinuationIfNeed(throwing: nil)
|
||||||
let channel = self.channel
|
let channel = self.channel
|
||||||
self.channel = nil
|
self.channel = nil
|
||||||
|
self.datagramSender.bind(nil, for: .v4)
|
||||||
try? channel?.close().wait()
|
try? channel?.close().wait()
|
||||||
try? self.group.syncShutdownGracefully()
|
try? self.group.syncShutdownGracefully()
|
||||||
|
|
||||||
|
|||||||
@ -15,18 +15,69 @@ 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 {
|
actor SDLUDPHoleService {
|
||||||
enum Event {
|
enum Event {
|
||||||
case ready(SocketAddress)
|
case ready(SocketAddress)
|
||||||
case natType(SDLNATProberActor.NatType)
|
case natType(SDLNATProberActor.NatType)
|
||||||
case packet(SocketAddress, SDLHoleMessage, source: SDLUDPHoleKind)
|
case packet(SocketAddress, SDLHoleControlMessage, source: SDLUDPHoleKind)
|
||||||
case closed(Error)
|
case closed(Error)
|
||||||
}
|
}
|
||||||
|
|
||||||
typealias EventHandler = @Sendable (Event) async -> Void
|
typealias EventHandler = @Sendable (Event) async -> Void
|
||||||
|
typealias DataHandler = @Sendable (SDLData) async -> Void
|
||||||
|
|
||||||
private let proberActor: SDLNATProberActor
|
private let proberActor: SDLNATProberActor
|
||||||
private let onEvent: EventHandler
|
private let onEvent: EventHandler
|
||||||
|
private let onData: DataHandler
|
||||||
|
private let datagramSender: SDLUDPHoleDatagramSender
|
||||||
|
|
||||||
private var udpHole: SDLUDPHole?
|
private var udpHole: SDLUDPHole?
|
||||||
private var udpHoleMonitorTask: Task<Void, Never>?
|
private var udpHoleMonitorTask: Task<Void, Never>?
|
||||||
@ -36,9 +87,16 @@ actor SDLUDPHoleService {
|
|||||||
private var udpHoleV6: SDLUDPHoleV6?
|
private var udpHoleV6: SDLUDPHoleV6?
|
||||||
private var udpHoleV6MonitorTask: Task<Void, Never>?
|
private var udpHoleV6MonitorTask: Task<Void, Never>?
|
||||||
|
|
||||||
init(proberActor: SDLNATProberActor, onEvent: @escaping EventHandler) {
|
init(
|
||||||
|
proberActor: SDLNATProberActor,
|
||||||
|
datagramSender: SDLUDPHoleDatagramSender,
|
||||||
|
onEvent: @escaping EventHandler,
|
||||||
|
onData: @escaping DataHandler
|
||||||
|
) {
|
||||||
self.proberActor = proberActor
|
self.proberActor = proberActor
|
||||||
|
self.datagramSender = datagramSender
|
||||||
self.onEvent = onEvent
|
self.onEvent = onEvent
|
||||||
|
self.onData = onData
|
||||||
}
|
}
|
||||||
|
|
||||||
func start(includeV6: Bool = false) {
|
func start(includeV6: Bool = false) {
|
||||||
@ -115,7 +173,7 @@ actor SDLUDPHoleService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private func runV4() async throws {
|
private func runV4() async throws {
|
||||||
let udpHole = try SDLUDPHole()
|
let udpHole = try SDLUDPHole(datagramSender: self.datagramSender)
|
||||||
let localAddress = try await udpHole.start()
|
let localAddress = try await udpHole.start()
|
||||||
self.udpHole = udpHole
|
self.udpHole = udpHole
|
||||||
self.localAddress = localAddress
|
self.localAddress = localAddress
|
||||||
@ -171,10 +229,10 @@ actor SDLUDPHoleService {
|
|||||||
case .stunProbeReply(let probeReply):
|
case .stunProbeReply(let probeReply):
|
||||||
await self.proberActor.handleProbeReply(localAddress: self.localAddress, reply: probeReply)
|
await self.proberActor.handleProbeReply(localAddress: self.localAddress, reply: probeReply)
|
||||||
default:
|
default:
|
||||||
await self.onEvent(.packet(remoteAddress, message, source: .v4))
|
await self.onEvent(.packet(remoteAddress, control, source: .v4))
|
||||||
}
|
}
|
||||||
case .data:
|
case .data(let data):
|
||||||
await self.onEvent(.packet(remoteAddress, message, source: .v4))
|
await self.onData(data)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -192,7 +250,7 @@ actor SDLUDPHoleService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private func runV6() async throws {
|
private func runV6() async throws {
|
||||||
let udpHoleV6 = try SDLUDPHoleV6()
|
let udpHoleV6 = try SDLUDPHoleV6(datagramSender: self.datagramSender)
|
||||||
let localAddress = try udpHoleV6.start()
|
let localAddress = try udpHoleV6.start()
|
||||||
self.udpHoleV6 = udpHoleV6
|
self.udpHoleV6 = udpHoleV6
|
||||||
|
|
||||||
@ -215,10 +273,16 @@ actor SDLUDPHoleService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let onEvent = self.onEvent
|
let onEvent = self.onEvent
|
||||||
|
let onData = self.onData
|
||||||
group.addTask {
|
group.addTask {
|
||||||
for await (remoteAddress, message) in udpHoleV6.messageStream {
|
for await (remoteAddress, message) in udpHoleV6.messageStream {
|
||||||
try Task.checkCancellation()
|
try Task.checkCancellation()
|
||||||
await onEvent(.packet(remoteAddress, message, source: .v6))
|
switch message {
|
||||||
|
case .control(let control):
|
||||||
|
await onEvent(.packet(remoteAddress, control, source: .v6))
|
||||||
|
case .data(let data):
|
||||||
|
await onData(data)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -242,22 +306,16 @@ actor SDLUDPHoleService {
|
|||||||
actor SDLUDPHoleServiceProxy {
|
actor SDLUDPHoleServiceProxy {
|
||||||
typealias ControlEventHandler = @Sendable (SDLUDPHoleService.Event) async -> Void
|
typealias ControlEventHandler = @Sendable (SDLUDPHoleService.Event) async -> Void
|
||||||
|
|
||||||
|
nonisolated let datagramSender = SDLUDPHoleDatagramSender()
|
||||||
private var udpHoleService: SDLUDPHoleService?
|
private var udpHoleService: SDLUDPHoleService?
|
||||||
private var packetInboundActor: PacketInboundActor?
|
|
||||||
private var onControlEvent: ControlEventHandler?
|
|
||||||
private var generation: UInt64 = 0
|
private var generation: UInt64 = 0
|
||||||
|
|
||||||
func bindInbound(_ packetInboundActor: PacketInboundActor, onControlEvent: @escaping ControlEventHandler) {
|
func makeEventHandler(onControlEvent: @escaping ControlEventHandler) -> SDLUDPHoleService.EventHandler {
|
||||||
self.packetInboundActor = packetInboundActor
|
|
||||||
self.onControlEvent = onControlEvent
|
|
||||||
}
|
|
||||||
|
|
||||||
func makeEventHandler() -> SDLUDPHoleService.EventHandler {
|
|
||||||
self.generation &+= 1
|
self.generation &+= 1
|
||||||
let generation = self.generation
|
let generation = self.generation
|
||||||
|
|
||||||
return { [weak self] event in
|
return { [weak self] event in
|
||||||
await self?.handleEvent(event, generation: generation)
|
await self?.handleEvent(event, generation: generation, onControlEvent: onControlEvent)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -283,21 +341,15 @@ actor SDLUDPHoleServiceProxy {
|
|||||||
await self.udpHoleService?.send(type: type, data: data, remoteAddress: remoteAddress)
|
await self.udpHoleService?.send(type: type, data: data, remoteAddress: remoteAddress)
|
||||||
}
|
}
|
||||||
|
|
||||||
private func handleEvent(_ event: SDLUDPHoleService.Event, generation: UInt64) async {
|
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 {
|
guard generation == self.generation else {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
switch event {
|
await onControlEvent(event)
|
||||||
case .packet(_, let message, _):
|
|
||||||
switch message {
|
|
||||||
case .data(let data):
|
|
||||||
await self.packetInboundActor?.handleData(data)
|
|
||||||
case .control:
|
|
||||||
await self.onControlEvent?(event)
|
|
||||||
}
|
|
||||||
case .ready, .natType, .closed:
|
|
||||||
await self.onControlEvent?(event)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -25,6 +25,7 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
|||||||
|
|
||||||
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
||||||
private var channel: Channel?
|
private var channel: Channel?
|
||||||
|
private let datagramSender: SDLUDPHoleDatagramSender
|
||||||
|
|
||||||
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
|
||||||
@ -34,7 +35,8 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
|||||||
private let eventContinuation: AsyncStream<HoleEvent>.Continuation
|
private let eventContinuation: AsyncStream<HoleEvent>.Continuation
|
||||||
|
|
||||||
// 启动函数
|
// 启动函数
|
||||||
init() throws {
|
init(datagramSender: SDLUDPHoleDatagramSender) throws {
|
||||||
|
self.datagramSender = datagramSender
|
||||||
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
|
||||||
@ -54,6 +56,7 @@ 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.datagramSender.bind(channel, for: .v6)
|
||||||
|
|
||||||
return channel.localAddress
|
return channel.localAddress
|
||||||
}
|
}
|
||||||
@ -121,6 +124,7 @@ final class SDLUDPHoleV6: ChannelInboundHandler {
|
|||||||
|
|
||||||
let channel = self.channel
|
let channel = self.channel
|
||||||
self.channel = nil
|
self.channel = nil
|
||||||
|
self.datagramSender.bind(nil, for: .v6)
|
||||||
try? channel?.close().wait()
|
try? channel?.close().wait()
|
||||||
|
|
||||||
try? self.group.syncShutdownGracefully()
|
try? self.group.syncShutdownGracefully()
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user