fix udpHole

This commit is contained in:
anlicheng 2026-05-22 17:04:05 +08:00
parent 28d23b3f8d
commit e07c81567a
4 changed files with 146 additions and 131 deletions

View File

@ -9,13 +9,6 @@ import NIOCore
import NIOPosix
import SwiftProtobuf
enum SDLUDPHoleError: Error {
case invalidLocalAddress
case closed
case errorCaught
case sendFaied(Error)
}
actor SDLUDPHole {
enum State {
case idle
@ -37,7 +30,7 @@ actor SDLUDPHole {
return localAddress
}
func messageStream() -> AsyncThrowingStream<(SocketAddress, SDLHoleMessage), Error> {
func messageStream() -> AsyncThrowingStream<SDLUDPHoleHandler.SDLHoleDatagram, Error> {
return self.udpHoleHandler.messageStream
}
@ -63,119 +56,3 @@ actor SDLUDPHole {
}
}
// sn-server
private final class SDLUDPHoleHandler: ChannelInboundHandler {
typealias InboundIn = AddressedEnvelope<ByteBuffer>
private let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
private var channel: Channel?
private let locker = NSLock()
public let messageStream: AsyncThrowingStream<(SocketAddress, SDLHoleMessage), Error>
private let messageContinuation: AsyncThrowingStream<(SocketAddress, SDLHoleMessage), Error>.Continuation
private var isMessageContinuationFinished: Bool = false
//
init() throws {
let (stream, continuation) = AsyncThrowingStream.makeStream(of: (SocketAddress, SDLHoleMessage).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((remoteAddress, message))
} else {
SDLLogger.log("[SDLUDPHole] decode message, get null", for: .debug)
}
} catch let err {
SDLLogger.log("[SDLUDPHole] decode message, get error: \(err)", for: .debug)
}
}
func channelInactive(context: ChannelHandlerContext) {
self.finishMessageContinuationIfNeed(throwing: .closed)
}
func errorCaught(context: ChannelHandlerContext, error: any Error) {
context.close(promise: nil)
self.finishMessageContinuationIfNeed(throwing: .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?.finishMessageContinuationIfNeed(throwing: .sendFaied(err))
}
}
func stop() {
self.finishMessageContinuationIfNeed(throwing: nil)
let channel = self.channel
self.channel = nil
try? channel?.close().wait()
try? self.group.syncShutdownGracefully()
SDLLogger.log("[SDLUDPHole] stopped", for: .debug)
}
private func finishMessageContinuationIfNeed(throwing error: SDLUDPHoleError?) {
locker.lock()
defer {
locker.unlock()
}
guard !self.isMessageContinuationFinished else {
return
}
self.isMessageContinuationFinished = true
self.messageContinuation.finish(throwing: error)
}
deinit {
SDLLogger.log("[SDLUDPHole] deinit", for: .debug)
}
}

View File

@ -0,0 +1,14 @@
//
// SDLUDPHoleError.swift
// punchnet
//
// Created by on 2026/5/22.
//
import Foundation
enum SDLUDPHoleError: Error {
case invalidLocalAddress
case closed
case errorCaught
case sendFaied(Error)
}

View File

@ -0,0 +1,128 @@
//
// SDLUDPHoleHandler.swift
// punchnet
//
// 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
private let locker = NSLock()
private var isStopped: Bool = false
//
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() {
locker.lock()
defer {
locker.unlock()
}
guard !self.isStopped else {
return
}
self.isStopped = true
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

@ -38,11 +38,7 @@ actor SDLUDPHoleService {
private var udpHoleV6: SDLUDPHoleV6?
private var udpHoleV6MonitorTask: Task<Void, Never>?
init(
proberActor: SDLNATProberActor,
onEvent: @escaping EventHandler,
onData: @escaping DataHandler
) {
init(proberActor: SDLNATProberActor, onEvent: @escaping EventHandler, onData: @escaping DataHandler) {
self.proberActor = proberActor
self.onEvent = onEvent
self.onData = onData
@ -140,9 +136,9 @@ actor SDLUDPHoleService {
do {
try await withTaskCancellationHandler {
for try await (remoteAddress, message) in await udpHole.messageStream() {
for try await datagram in await udpHole.messageStream() {
try Task.checkCancellation()
try await self.handleV4Message(remoteAddress: remoteAddress, message: message)
try await self.handleV4Message(remoteAddress: datagram.remoteAddress, message: datagram.message)
}
} onCancel: {
Task {