拆分udpHoleV6

This commit is contained in:
anlicheng 2026-05-27 21:24:22 +08:00
parent 3f0e0a86a9
commit 0c31f1425a
12 changed files with 330 additions and 248 deletions

View File

@ -43,6 +43,7 @@ actor SDLContextActor {
private var dnsService: DNSService? private var dnsService: DNSService?
private let superService: SDLSuperService private let superService: SDLSuperService
private let udpHoleService: SDLUDPHoleService private let udpHoleService: SDLUDPHoleService
private let udpHoleV6Service: SDLUDPHoleV6Service
private let packetOutboundActor: PacketOutboundActor private let packetOutboundActor: PacketOutboundActor
private let packetInboundActor: PacketInboundActor private let packetInboundActor: PacketInboundActor
private let tunNetworkManager: SDLTunNetworkManager private let tunNetworkManager: SDLTunNetworkManager
@ -83,6 +84,7 @@ actor SDLContextActor {
let policyService = PolicyService(identityId: config.identityId, acl: config.acl) let policyService = PolicyService(identityId: config.identityId, acl: config.acl)
let superService = SDLSuperService(serverEndpoint: config.serverEndpoint) let superService = SDLSuperService(serverEndpoint: config.serverEndpoint)
let udpHoleService = SDLUDPHoleService(proberActor: proberActor) let udpHoleService = SDLUDPHoleService(proberActor: proberActor)
let udpHoleV6Service = SDLUDPHoleV6Service()
let tunNetworkManager = SDLTunNetworkManager(provider: provider) let tunNetworkManager = SDLTunNetworkManager(provider: provider)
let ipv6AssistPair = AsyncStream.makeStream(of: Optional<SDLV6Info>.self, bufferingPolicy: .bufferingNewest(1)) let ipv6AssistPair = AsyncStream.makeStream(of: Optional<SDLV6Info>.self, bufferingPolicy: .bufferingNewest(1))
let packetOutboundActor = PacketOutboundActor( let packetOutboundActor = PacketOutboundActor(
@ -95,6 +97,7 @@ actor SDLContextActor {
policyService: policyService, policyService: policyService,
superService: superService, superService: superService,
udpHoleService: udpHoleService, udpHoleService: udpHoleService,
udpHoleV6Service: udpHoleV6Service,
flowTracer: flowTracer flowTracer: flowTracer
) )
let packetInboundActor = PacketInboundActor( let packetInboundActor = PacketInboundActor(
@ -123,6 +126,7 @@ actor SDLContextActor {
self.policyService = policyService self.policyService = policyService
self.superService = superService self.superService = superService
self.udpHoleService = udpHoleService self.udpHoleService = udpHoleService
self.udpHoleV6Service = udpHoleV6Service
self.packetOutboundActor = packetOutboundActor self.packetOutboundActor = packetOutboundActor
self.packetInboundActor = packetInboundActor self.packetInboundActor = packetInboundActor
self.tunNetworkManager = tunNetworkManager self.tunNetworkManager = tunNetworkManager
@ -198,9 +202,18 @@ actor SDLContextActor {
await packetInboundActor.handleData(data) await packetInboundActor.handleData(data)
} }
) )
await self.udpHoleV6Service.updateHandlers(
onEvent: { [weak self] event in
await self?.handleUDPHoleControlEvent(event)
},
onData: { data in
await packetInboundActor.handleData(data)
}
)
let superService = self.superService let superService = self.superService
let udpHoleService = self.udpHoleService let udpHoleService = self.udpHoleService
let udpHoleV6Service = self.udpHoleV6Service
let packetOutboundActor = self.packetOutboundActor let packetOutboundActor = self.packetOutboundActor
let policyService = self.policyService let policyService = self.policyService
let puncherActor = self.puncherActor let puncherActor = self.puncherActor
@ -220,7 +233,13 @@ actor SDLContextActor {
group.addTask { group.addTask {
try await Self.runRestarting(name: "udpHoleService") { try await Self.runRestarting(name: "udpHoleService") {
try await udpHoleService.run(includeV6: false) try await udpHoleService.run()
}
}
group.addTask {
try await Self.runRestarting(name: "udpHoleV6Service") {
try await udpHoleV6Service.run()
} }
} }
@ -357,6 +376,7 @@ actor SDLContextActor {
await self.policyService.clear() await self.policyService.clear()
await self.udpHoleService.stop() await self.udpHoleService.stop()
await self.udpHoleV6Service.stop()
let dnsService = self.dnsService let dnsService = self.dnsService
self.dnsService = nil self.dnsService = nil
@ -419,7 +439,14 @@ extension SDLContextActor {
} }
private func sendPacket(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async { private func sendPacket(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async {
await self.udpHoleService.send(type: type, data: data, remoteAddress: remoteAddress) switch remoteAddress {
case .v4:
await self.udpHoleService.send(type: type, data: data, remoteAddress: remoteAddress)
case .v6:
await self.udpHoleV6Service.send(type: type, data: data, remoteAddress: remoteAddress)
default:
SDLLogger.log("[SDLContext] unsupported socket family: \(remoteAddress)", for: .debug)
}
} }
} }
@ -453,7 +480,7 @@ extension SDLContextActor {
SDLLogger.log("[SDLContext] peer message: \(peerInfo)") SDLLogger.log("[SDLContext] peer message: \(peerInfo)")
let packets = await self.puncherActor.makeRegisterPackets(peerInfo: peerInfo) let packets = await self.puncherActor.makeRegisterPackets(peerInfo: peerInfo)
for packet in packets { for packet in packets {
await self.udpHoleService.send(type: .register, data: packet.data, remoteAddress: packet.remoteAddress) await self.sendPacket(type: .register, data: packet.data, remoteAddress: packet.remoteAddress)
} }
case .event(let event): case .event(let event):
await self.handleEvent(event: event) await self.handleEvent(event: event)

View File

@ -25,6 +25,7 @@ actor PacketOutboundActor {
private let policyService: PolicyService private let policyService: PolicyService
private let superService: SDLSuperService private let superService: SDLSuperService
private let udpHoleService: SDLUDPHoleService private let udpHoleService: SDLUDPHoleService
private let udpHoleV6Service: SDLUDPHoleV6Service
private let flowTracer: SDLFlowTracer private let flowTracer: SDLFlowTracer
private var networkAddress: SDLConfiguration.NetworkAddress private var networkAddress: SDLConfiguration.NetworkAddress
@ -43,6 +44,7 @@ actor PacketOutboundActor {
policyService: PolicyService, policyService: PolicyService,
superService: SDLSuperService, superService: SDLSuperService,
udpHoleService: SDLUDPHoleService, udpHoleService: SDLUDPHoleService,
udpHoleV6Service: SDLUDPHoleV6Service,
flowTracer: SDLFlowTracer) { flowTracer: SDLFlowTracer) {
self.provider = provider self.provider = provider
self.networkAddress = config.networkAddress self.networkAddress = config.networkAddress
@ -56,6 +58,7 @@ actor PacketOutboundActor {
self.policyService = policyService self.policyService = policyService
self.superService = superService self.superService = superService
self.udpHoleService = udpHoleService self.udpHoleService = udpHoleService
self.udpHoleV6Service = udpHoleV6Service
self.flowTracer = flowTracer self.flowTracer = flowTracer
} }
@ -233,6 +236,13 @@ actor PacketOutboundActor {
} }
private func sendPacket(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async { private func sendPacket(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async {
await self.udpHoleService.send(type: type, data: data, remoteAddress: remoteAddress) switch remoteAddress {
case .v4:
await self.udpHoleService.send(type: type, data: data, remoteAddress: remoteAddress)
case .v6:
await self.udpHoleV6Service.send(type: type, data: data, remoteAddress: remoteAddress)
default:
SDLLogger.log("[PacketOutboundActor] unsupported socket family: \(remoteAddress)", for: .debug)
}
} }
} }

View File

@ -175,19 +175,19 @@ actor SDLNATProberActor {
guard !Task.isCancelled else { guard !Task.isCancelled else {
return return
} }
await udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 1, attr: .none), remoteAddress: addressArray[0][0]) udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 1, attr: .none), remoteAddress: addressArray[0][0])
guard !Task.isCancelled else { guard !Task.isCancelled else {
return return
} }
await udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 2, attr: .none), remoteAddress: addressArray[1][1]) udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 2, attr: .none), remoteAddress: addressArray[1][1])
guard !Task.isCancelled else { guard !Task.isCancelled else {
return return
} }
await udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 3, attr: .peer), remoteAddress: addressArray[0][0]) udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 3, attr: .peer), remoteAddress: addressArray[0][0])
guard !Task.isCancelled else { guard !Task.isCancelled else {
return return
} }
await udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 4, attr: .port), remoteAddress: addressArray[0][0]) udpHole.send(type: .stunProbe, data: makeProbePacket(cookieId: cookie, step: 4, attr: .port), remoteAddress: addressArray[0][0])
} }
private func makeProbePacket(cookieId: UInt32, step: UInt32, attr: SDLProbeAttr) -> Data { private func makeProbePacket(cookieId: UInt32, step: UInt32, attr: SDLProbeAttr) -> Data {

View File

@ -7,52 +7,129 @@
import Foundation import Foundation
import NIOCore import NIOCore
import NIOPosix import NIOPosix
import SwiftProtobuf
actor SDLUDPHole { // sn-server
final class SDLUDPHole: ChannelInboundHandler {
typealias InboundIn = AddressedEnvelope<ByteBuffer>
struct SDLHoleDatagram {
let remoteAddress: SocketAddress
let message: SDLHoleMessage
}
enum State { enum State {
case idle case idle
case running case running
case stopped case stopped
} }
private var state: State = .idle private var state: State = .idle
private let udpHoleHandler: SDLUDPHoleHandler private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
private var channel: Channel?
let messageStream: AsyncThrowingStream<SDLHoleDatagram, Error>
private let messageContinuation: AsyncThrowingStream<SDLHoleDatagram, Error>.Continuation
init() throws { init() throws {
self.udpHoleHandler = try SDLUDPHoleHandler() let (stream, continuation) = AsyncThrowingStream.makeStream(of: SDLHoleDatagram.self, bufferingPolicy: .bufferingNewest(2048))
self.messageStream = stream
self.messageContinuation = continuation
} }
func start() async throws -> SocketAddress { func start() throws -> SocketAddress {
let localAddress = try self.udpHoleHandler.start() guard self.state == .idle else {
guard let localAddress = self.channel?.localAddress else {
throw SDLUDPHoleError.invalidLocalAddress
}
return localAddress
}
let bootstrap = DatagramBootstrap(group: group)
.channelOption(ChannelOptions.socketOption(.so_reuseaddr), value: 1)
.channelInitializer { channel in
channel.pipeline.addHandler(self)
}
// IPv4IPv4
let channel = try bootstrap.bind(host: "0.0.0.0", port: 0).wait()
guard let localAddress = channel.localAddress else {
throw SDLUDPHoleError.invalidLocalAddress
}
self.channel = channel
self.state = .running self.state = .running
return localAddress return localAddress
} }
func messageStream() -> AsyncThrowingStream<SDLUDPHoleHandler.SDLHoleDatagram, Error> { // --MARK: ChannelInboundHandler delegate
return self.udpHoleHandler.messageStream
func channelRead(context: ChannelHandlerContext, data: NIOAny) {
let envelope = unwrapInboundIn(data)
var buffer = envelope.data
let remoteAddress = envelope.remoteAddress
do {
if let message = try SDLHoleMessage.decode(buffer: &buffer) {
self.messageContinuation.yield(SDLHoleDatagram(remoteAddress: remoteAddress, message: message))
} else {
SDLLogger.log("[SDLUDPHole] decode message, get null", for: .debug)
}
} catch let err {
SDLLogger.log("[SDLUDPHole] decode message, get error: \(err)", for: .debug)
self.messageContinuation.finish(throwing: err)
}
} }
func channelInactive(context: ChannelHandlerContext) {
self.messageContinuation.finish(throwing: SDLUDPHoleError.closed)
}
func errorCaught(context: ChannelHandlerContext, error: any Error) {
context.close(promise: nil)
self.messageContinuation.finish(throwing: SDLUDPHoleError.errorCaught)
}
// MARK:
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) { func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) {
guard self.state == .running else { guard self.state == .running, let channel = self.channel else {
return return
} }
self.udpHoleHandler.send(type: type, data: data, remoteAddress: remoteAddress) var buffer = channel.allocator.buffer(capacity: data.count + 1)
buffer.writeBytes([type.rawValue])
buffer.writeBytes(data)
let envelope = AddressedEnvelope<ByteBuffer>(remoteAddress: remoteAddress, data: buffer)
let promise = channel.eventLoop.makePromise(of: Void.self)
channel.eventLoop.execute {
channel.writeAndFlush(envelope, promise: promise)
}
promise.futureResult.whenFailure { [weak self] err in
self?.messageContinuation.finish(throwing: SDLUDPHoleError.sendFaied(err))
}
} }
func stop() async { func stop() {
guard self.state != .stopped else { guard self.state != .stopped else {
return return
} }
self.state = .stopped self.state = .stopped
self.udpHoleHandler.stop() self.messageContinuation.finish()
let channel = self.channel
self.channel = nil
try? channel?.close().wait()
try? self.group.syncShutdownGracefully()
SDLLogger.log("[SDLUDPHole] stopped", for: .debug)
} }
deinit { deinit {
SDLLogger.log("[SDLUDPHoleActor] deinit", for: .debug) SDLLogger.log("[SDLUDPHole] deinit", for: .debug)
} }
} }

View File

@ -1,115 +0,0 @@
//
// SDLUDPHoleHandler.swift
// punchnet
// SDLUDPHoleactorSDLUDPHoleHandler
// Created by on 2026/5/22.
//
import Foundation
import NIOCore
import NIOPosix
// sn-server
final class SDLUDPHoleHandler: ChannelInboundHandler {
typealias InboundIn = AddressedEnvelope<ByteBuffer>
struct SDLHoleDatagram {
let remoteAddress: SocketAddress
let message: SDLHoleMessage
}
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
private var channel: Channel?
public let messageStream: AsyncThrowingStream<SDLHoleDatagram, Error>
private let messageContinuation: AsyncThrowingStream<SDLHoleDatagram, Error>.Continuation
//
init() throws {
let (stream, continuation) = AsyncThrowingStream.makeStream(of: SDLHoleDatagram.self, bufferingPolicy: .bufferingNewest(2048))
self.messageStream = stream
self.messageContinuation = continuation
}
func start() throws -> SocketAddress {
let bootstrap = DatagramBootstrap(group: group)
.channelOption(ChannelOptions.socketOption(.so_reuseaddr), value: 1)
.channelInitializer { channel in
channel.pipeline.addHandler(self)
}
// IPv4IPv4
let channel = try bootstrap.bind(host: "0.0.0.0", port: 0).wait()
guard let localAddress = channel.localAddress else {
throw SDLUDPHoleError.invalidLocalAddress
}
self.channel = channel
return localAddress
}
// --MARK: ChannelInboundHandler delegate
func channelRead(context: ChannelHandlerContext, data: NIOAny) {
let envelope = unwrapInboundIn(data)
var buffer = envelope.data
let remoteAddress = envelope.remoteAddress
do {
if let message = try SDLHoleMessage.decode(buffer: &buffer) {
self.messageContinuation.yield(SDLHoleDatagram(remoteAddress: remoteAddress, message: message))
} else {
SDLLogger.log("[SDLUDPHole] decode message, get null", for: .debug)
}
} catch let err {
SDLLogger.log("[SDLUDPHole] decode message, get error: \(err)", for: .debug)
self.messageContinuation.finish(throwing: err)
}
}
func channelInactive(context: ChannelHandlerContext) {
self.messageContinuation.finish(throwing: SDLUDPHoleError.closed)
}
func errorCaught(context: ChannelHandlerContext, error: any Error) {
context.close(promise: nil)
self.messageContinuation.finish(throwing: SDLUDPHoleError.errorCaught)
}
// MARK:
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) {
guard let channel = self.channel else {
return
}
var buffer = channel.allocator.buffer(capacity: data.count + 1)
buffer.writeBytes([type.rawValue])
buffer.writeBytes(data)
let envelope = AddressedEnvelope<ByteBuffer>(remoteAddress: remoteAddress, data: buffer)
let promise = channel.eventLoop.makePromise(of: Void.self)
channel.eventLoop.execute {
channel.writeAndFlush(envelope, promise: promise)
}
promise.futureResult.whenFailure { [weak self] err in
self?.messageContinuation.finish(throwing: SDLUDPHoleError.sendFaied(err))
}
}
func stop() {
self.messageContinuation.finish()
let channel = self.channel
self.channel = nil
try? channel?.close().wait()
try? self.group.syncShutdownGracefully()
SDLLogger.log("[SDLUDPHole] stopped", for: .debug)
}
deinit {
SDLLogger.log("[SDLUDPHole] deinit", for: .debug)
}
}

View File

@ -42,11 +42,10 @@ actor SDLUDPHoleService {
self.onData = onData self.onData = onData
} }
func run(includeV6: Bool = false) async throws { func run() async throws {
let generation = self.nextGeneration() let generation = self.nextGeneration()
let session = SDLUDPHoleSession( let session = SDLUDPHoleSession(
proberActor: self.proberActor, proberActor: self.proberActor,
includeV6: includeV6,
onEvent: { [weak self] event in onEvent: { [weak self] event in
await self?.handleEvent(event, generation: generation) await self?.handleEvent(event, generation: generation)
}, },

View File

@ -9,45 +9,25 @@ import NIOCore
actor SDLUDPHoleSession { actor SDLUDPHoleSession {
private let proberActor: SDLNATProberActor private let proberActor: SDLNATProberActor
private let includeV6: Bool
private let onEvent: SDLUDPHoleService.EventHandler private let onEvent: SDLUDPHoleService.EventHandler
private let onData: SDLUDPHoleService.DataHandler private let onData: SDLUDPHoleService.DataHandler
private var udpHole: SDLUDPHole? private var udpHole: SDLUDPHole?
private var udpHoleV6: SDLUDPHoleV6?
private var localAddress: SocketAddress? private var localAddress: SocketAddress?
init( init(
proberActor: SDLNATProberActor, proberActor: SDLNATProberActor,
includeV6: Bool,
onEvent: @escaping SDLUDPHoleService.EventHandler, onEvent: @escaping SDLUDPHoleService.EventHandler,
onData: @escaping SDLUDPHoleService.DataHandler onData: @escaping SDLUDPHoleService.DataHandler
) { ) {
self.proberActor = proberActor self.proberActor = proberActor
self.includeV6 = includeV6
self.onEvent = onEvent self.onEvent = onEvent
self.onData = onData self.onData = onData
} }
func run() async throws { func run() async throws {
do { do {
try await withThrowingTaskGroup(of: Void.self) { group in try await self.runV4()
defer {
group.cancelAll()
}
group.addTask {
try await self.runV4()
}
if self.includeV6 {
group.addTask {
try await self.runV6()
}
}
try await group.waitForAll()
}
await self.stop() await self.stop()
} catch { } catch {
await self.stop() await self.stop()
@ -59,39 +39,28 @@ actor SDLUDPHoleSession {
let udpHole = self.udpHole let udpHole = self.udpHole
self.udpHole = nil self.udpHole = nil
self.localAddress = nil self.localAddress = nil
await udpHole?.stop() udpHole?.stop()
let udpHoleV6 = self.udpHoleV6
self.udpHoleV6 = nil
udpHoleV6?.stop()
await self.proberActor.cancelAll() await self.proberActor.cancelAll()
} }
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async { func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async {
switch remoteAddress { guard case .v4 = remoteAddress else {
case .v4:
guard let udpHole else {
SDLLogger.log("[SDLUDPHoleSession] udpHole is nil for remoteAddress: \(remoteAddress)", for: .debug)
return
}
await udpHole.send(type: type, data: data, remoteAddress: remoteAddress)
case .v6:
guard let udpHoleV6 else {
SDLLogger.log("[SDLUDPHoleSession] udpHoleV6 is nil for remoteAddress: \(remoteAddress)", for: .debug)
return
}
udpHoleV6.send(type: type, data: data, remoteAddress: remoteAddress)
default:
SDLLogger.log("[SDLUDPHoleSession] unsupported socket family: \(remoteAddress)", for: .debug) SDLLogger.log("[SDLUDPHoleSession] unsupported socket family: \(remoteAddress)", for: .debug)
return
} }
guard let udpHole else {
SDLLogger.log("[SDLUDPHoleSession] udpHole is nil for remoteAddress: \(remoteAddress)", for: .debug)
return
}
udpHole.send(type: type, data: data, remoteAddress: remoteAddress)
} }
private func runV4() async throws { private func runV4() async throws {
let udpHole = try SDLUDPHole() let udpHole = try SDLUDPHole()
let localAddress = try await udpHole.start() let localAddress = try udpHole.start()
self.udpHole = udpHole self.udpHole = udpHole
self.localAddress = localAddress self.localAddress = localAddress
@ -115,7 +84,7 @@ actor SDLUDPHoleSession {
try await group.waitForAll() try await group.waitForAll()
} }
} catch { } catch {
await udpHole.stop() udpHole.stop()
if self.udpHole === udpHole { if self.udpHole === udpHole {
self.udpHole = nil self.udpHole = nil
self.localAddress = nil self.localAddress = nil
@ -125,7 +94,7 @@ actor SDLUDPHoleSession {
} }
private func readV4Loop(udpHole: SDLUDPHole) async throws { private func readV4Loop(udpHole: SDLUDPHole) async throws {
for try await datagram in await udpHole.messageStream() { for try await datagram in udpHole.messageStream {
try Task.checkCancellation() try Task.checkCancellation()
try await self.handleV4Message(remoteAddress: datagram.remoteAddress, message: datagram.message) try await self.handleV4Message(remoteAddress: datagram.remoteAddress, message: datagram.message)
} }
@ -157,58 +126,4 @@ actor SDLUDPHoleSession {
await self.onData(data) await self.onData(data)
} }
} }
private func runV6() async throws {
let udpHoleV6 = try SDLUDPHoleV6()
let localAddress = try udpHoleV6.start()
self.udpHoleV6 = udpHoleV6
if let localAddress {
SDLLogger.log("[SDLUDPHoleSession] udpHoleV6 started, on address: \(localAddress)")
} else {
SDLLogger.log("[SDLUDPHoleSession] udpHoleV6 started, no local address")
}
do {
try await withThrowingTaskGroup(of: Void.self) { group in
defer {
group.cancelAll()
}
let onEvent = self.onEvent
let onData = self.onData
group.addTask {
for await (remoteAddress, message) in udpHoleV6.messageStream {
try Task.checkCancellation()
switch message {
case .control(let control):
await onEvent(.packet(remoteAddress, control, source: .v6))
case .data(let data):
await onData(data)
}
}
}
group.addTask {
for await event in udpHoleV6.eventStream {
try Task.checkCancellation()
switch event {
case .ready:
SDLLogger.log("[SDLUDPHoleSession] udpHoleV6 ready")
case .closed, .errorCaught:
throw SDLContextError.udpHoleClosed
}
}
}
_ = try await group.next()
}
} catch {
udpHoleV6.stop()
if self.udpHoleV6 === udpHoleV6 {
self.udpHoleV6 = nil
}
throw error
}
}
} }

View File

@ -0,0 +1,78 @@
import Foundation
import NIOCore
actor SDLUDPHoleV6Service {
typealias EventHandler = SDLUDPHoleService.EventHandler
typealias DataHandler = SDLUDPHoleService.DataHandler
private var onEvent: EventHandler = { _ in }
private var onData: DataHandler = { _ in }
private var currentSession: SDLUDPHoleV6Session?
private var generation: UInt64 = 0
func updateHandlers(onEvent: @escaping EventHandler, onData: @escaping DataHandler) {
self.onEvent = onEvent
self.onData = onData
}
func run() async throws {
let generation = self.nextGeneration()
let session = SDLUDPHoleV6Session(
onEvent: { [weak self] event in
await self?.handleEvent(event, generation: generation)
},
onData: self.onData
)
self.currentSession = session
do {
try await session.run()
self.clearCurrent(session, generation: generation)
} catch is CancellationError {
self.clearCurrent(session, generation: generation)
await session.stop()
throw CancellationError()
} catch {
self.clearCurrent(session, generation: generation)
await session.stop()
throw error
}
}
func stop() async {
self.generation &+= 1
let session = self.currentSession
self.currentSession = nil
await session?.stop()
}
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async {
await self.currentSession?.send(type: type, data: data, remoteAddress: remoteAddress)
}
private func nextGeneration() -> UInt64 {
self.generation &+= 1
return self.generation
}
private func clearCurrent(_ session: SDLUDPHoleV6Session, generation: UInt64) {
guard self.generation == generation else {
return
}
if self.currentSession === session {
self.currentSession = nil
}
}
private func handleEvent(_ event: SDLUDPHoleService.Event, generation: UInt64) async {
guard self.generation == generation else {
return
}
await self.onEvent(event)
}
}

View File

@ -0,0 +1,91 @@
import Foundation
import NIOCore
actor SDLUDPHoleV6Session {
private let onEvent: SDLUDPHoleService.EventHandler
private let onData: SDLUDPHoleService.DataHandler
private var udpHoleV6: SDLUDPHoleV6?
init(
onEvent: @escaping SDLUDPHoleService.EventHandler,
onData: @escaping SDLUDPHoleService.DataHandler
) {
self.onEvent = onEvent
self.onData = onData
}
func run() async throws {
let udpHoleV6 = try SDLUDPHoleV6()
let localAddress = try udpHoleV6.start()
self.udpHoleV6 = udpHoleV6
if let localAddress {
SDLLogger.log("[SDLUDPHoleV6Session] udpHoleV6 started, on address: \(localAddress)")
} else {
SDLLogger.log("[SDLUDPHoleV6Session] udpHoleV6 started, no local address")
}
do {
try await withThrowingTaskGroup(of: Void.self) { group in
defer {
group.cancelAll()
}
let onEvent = self.onEvent
let onData = self.onData
group.addTask {
for await (remoteAddress, message) in udpHoleV6.messageStream {
try Task.checkCancellation()
switch message {
case .control(let control):
await onEvent(.packet(remoteAddress, control, source: .v6))
case .data(let data):
await onData(data)
}
}
}
group.addTask {
for await event in udpHoleV6.eventStream {
try Task.checkCancellation()
switch event {
case .ready:
SDLLogger.log("[SDLUDPHoleV6Session] udpHoleV6 ready")
case .closed, .errorCaught:
throw SDLContextError.udpHoleClosed
}
}
}
_ = try await group.next()
}
await self.stop()
} catch {
await self.stop()
throw error
}
}
func stop() async {
let udpHoleV6 = self.udpHoleV6
self.udpHoleV6 = nil
udpHoleV6?.stop()
}
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) {
guard case .v6 = remoteAddress else {
SDLLogger.log("[SDLUDPHoleV6Session] unsupported socket family: \(remoteAddress)", for: .debug)
return
}
guard let udpHoleV6 else {
SDLLogger.log("[SDLUDPHoleV6Session] udpHoleV6 is nil for remoteAddress: \(remoteAddress)", for: .debug)
return
}
udpHoleV6.send(type: type, data: data, remoteAddress: remoteAddress)
}
}