fix inbound
This commit is contained in:
parent
e6c79d7519
commit
304324e0d5
@ -67,6 +67,7 @@ actor SDLContextActor {
|
||||
private let superServiceProxy: SDLSuperServiceProxy
|
||||
private let udpHoleServiceProxy: SDLUDPHoleServiceProxy
|
||||
private let packetOutboundActor: PacketOutboundActor
|
||||
private let packetInboundActor: PacketInboundActor
|
||||
|
||||
private let publicDnsServers = ["223.5.5.5", "119.29.29.29"]
|
||||
|
||||
@ -104,6 +105,28 @@ actor SDLContextActor {
|
||||
let policyService = PolicyService(identityId: config.identityId)
|
||||
let superServiceProxy = SDLSuperServiceProxy()
|
||||
let udpHoleServiceProxy = SDLUDPHoleServiceProxy()
|
||||
let packetOutboundActor = PacketOutboundActor(
|
||||
provider: provider,
|
||||
config: config,
|
||||
dataCipher: nil,
|
||||
sessionManager: sessionManager,
|
||||
arpServer: arpServer,
|
||||
puncherActor: puncherActor,
|
||||
policyService: policyService,
|
||||
superServiceProxy: superServiceProxy,
|
||||
udpHoleServiceProxy: udpHoleServiceProxy,
|
||||
flowTracer: flowTracer
|
||||
)
|
||||
let packetInboundActor = PacketInboundActor(
|
||||
provider: provider,
|
||||
config: config,
|
||||
dataCipher: nil,
|
||||
policyService: policyService,
|
||||
packetOutboundActor: packetOutboundActor,
|
||||
arpServer: arpServer,
|
||||
superServiceProxy: superServiceProxy,
|
||||
flowTracer: flowTracer
|
||||
)
|
||||
|
||||
self.provider = provider
|
||||
self.config = config
|
||||
@ -120,18 +143,8 @@ actor SDLContextActor {
|
||||
self.policyService = policyService
|
||||
self.superServiceProxy = superServiceProxy
|
||||
self.udpHoleServiceProxy = udpHoleServiceProxy
|
||||
self.packetOutboundActor = PacketOutboundActor(
|
||||
provider: provider,
|
||||
config: config,
|
||||
dataCipher: nil,
|
||||
sessionManager: sessionManager,
|
||||
arpServer: arpServer,
|
||||
puncherActor: puncherActor,
|
||||
policyService: policyService,
|
||||
superServiceProxy: superServiceProxy,
|
||||
udpHoleServiceProxy: udpHoleServiceProxy,
|
||||
flowTracer: flowTracer
|
||||
)
|
||||
self.packetOutboundActor = packetOutboundActor
|
||||
self.packetInboundActor = packetInboundActor
|
||||
}
|
||||
|
||||
public func start() async {
|
||||
@ -148,9 +161,12 @@ actor SDLContextActor {
|
||||
await self.packetOutboundActor.updateDNSService(dnsService)
|
||||
await dnsService.start()
|
||||
|
||||
let udpHoleService = SDLUDPHoleService(proberActor: self.proberActor) { [weak self] event in
|
||||
await self?.handleUDPHoleEvent(event)
|
||||
await self.udpHoleServiceProxy.bindInbound(self.packetInboundActor) { [weak self] event in
|
||||
await self?.handleUDPHoleControlEvent(event)
|
||||
}
|
||||
|
||||
let udpHoleEventHandler = await self.udpHoleServiceProxy.makeEventHandler()
|
||||
let udpHoleService = SDLUDPHoleService(proberActor: self.proberActor, onEvent: udpHoleEventHandler)
|
||||
await self.udpHoleServiceProxy.replace(udpHoleService)
|
||||
await udpHoleService.start(includeV6: false)
|
||||
|
||||
@ -189,6 +205,7 @@ actor SDLContextActor {
|
||||
self.dataCipher = nil
|
||||
self.natType = .blocked
|
||||
await self.packetOutboundActor.updateRuntime(config: self.config, dataCipher: nil)
|
||||
await self.packetInboundActor.updateRuntime(config: self.config, dataCipher: nil)
|
||||
|
||||
await self.ipv6AssistClient?.stop()
|
||||
self.ipv6AssistClient = nil
|
||||
@ -315,6 +332,7 @@ extension SDLContextActor {
|
||||
}
|
||||
|
||||
await self.packetOutboundActor.updateRuntime(config: self.config, dataCipher: self.dataCipher)
|
||||
await self.packetInboundActor.updateRuntime(config: self.config, dataCipher: self.dataCipher)
|
||||
SDLLogger.log("[SDLContext] registerSuperAck, use algorithm \(algorithm), key len: \(key.count)")
|
||||
// 服务器分配的tun网卡信息
|
||||
do {
|
||||
@ -426,7 +444,7 @@ extension SDLContextActor {
|
||||
|
||||
// MARK: 处理从Hole收到的数据
|
||||
extension SDLContextActor {
|
||||
private func handleUDPHoleEvent(_ event: SDLUDPHoleService.Event) async {
|
||||
private func handleUDPHoleControlEvent(_ event: SDLUDPHoleService.Event) async {
|
||||
switch event {
|
||||
case .ready(let localAddress):
|
||||
SDLLogger.log("[SDLContext] udpHole ready: \(localAddress)")
|
||||
@ -451,8 +469,8 @@ extension SDLContextActor {
|
||||
case .registerAck(let registerAck):
|
||||
await self.handleRegisterAck(remoteAddress: remoteAddress, registerAck: registerAck, source: source)
|
||||
}
|
||||
case .data(let data):
|
||||
try? await self.handleHoleData(data: data)
|
||||
case .data:
|
||||
SDLLogger.log("[SDLContext] unexpected data packet in control path", for: .debug)
|
||||
}
|
||||
}
|
||||
|
||||
@ -494,42 +512,6 @@ extension SDLContextActor {
|
||||
}
|
||||
}
|
||||
|
||||
private func makeHoleDataProcessor() -> SDLHoleDataProcessor {
|
||||
return .init(
|
||||
networkAddress: self.config.networkAddress,
|
||||
dataCipher: self.dataCipher,
|
||||
policyService: self.policyService)
|
||||
}
|
||||
|
||||
private func handleHoleData(data: SDLData) async throws {
|
||||
let processor = self.makeHoleDataProcessor()
|
||||
guard let plan = try await processor.makeProcessingPlan(data: data) else {
|
||||
return
|
||||
}
|
||||
|
||||
self.flowTracer.inc(num: plan.inboundBytes, type: .inbound)
|
||||
|
||||
switch plan.action {
|
||||
case .sendARPReply(let dstMac, let responseData):
|
||||
SDLLogger.log("[SDLContext] get arp request packet")
|
||||
await self.packetOutboundActor.routeLayerPacket(dstMac: dstMac, type: .arp, data: responseData)
|
||||
case .appendARP(let ip, let mac):
|
||||
SDLLogger.log("[SDLContext] get arp response packet")
|
||||
await self.arpServer.append(ip: ip, mac: mac)
|
||||
case .writeToTun(let packetData, let identityID):
|
||||
let packet = NEPacket(data: packetData, protocolFamily: 2)
|
||||
self.provider.packetFlow.writePacketObjects([packet])
|
||||
SDLLogger.log("[SDLContext] hole identity: \(identityID), allow, data count: \(packetData.count)", for: .trace)
|
||||
case .requestPolicy(let srcIdentityID):
|
||||
SDLLogger.log("[SDLContext] not found identity: \(srcIdentityID) ruleMap", for: .debug)
|
||||
if let queryData = await self.policyService.identifyStore.makePolicyRequest(srcIdentityId: srcIdentityID, dstIdentityId: self.config.identityId) {
|
||||
await self.superServiceProxy.send(type: .policyRequest, data: queryData)
|
||||
}
|
||||
case .none:
|
||||
()
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
// MARK: 和Stun相关的心跳机制
|
||||
@ -606,6 +588,7 @@ extension SDLContextActor {
|
||||
self.config.exitNode = nil
|
||||
}
|
||||
await self.packetOutboundActor.updateRuntime(config: self.config, dataCipher: self.dataCipher)
|
||||
await self.packetInboundActor.updateRuntime(config: self.config, dataCipher: self.dataCipher)
|
||||
try await self.setNetworkSettings(config: self.config, dnsServer: DNSHelper.dnsServer)
|
||||
}
|
||||
|
||||
|
||||
81
Tun/Punchnet/Inbound/PacketInboundActor.swift
Normal file
81
Tun/Punchnet/Inbound/PacketInboundActor.swift
Normal file
@ -0,0 +1,81 @@
|
||||
//
|
||||
// PacketInboundActor.swift
|
||||
// Tun
|
||||
//
|
||||
// Created by Codex on 2026/5/20.
|
||||
//
|
||||
|
||||
import Foundation
|
||||
import NetworkExtension
|
||||
|
||||
actor PacketInboundActor {
|
||||
private let provider: NEPacketTunnelProvider
|
||||
private let policyService: PolicyService
|
||||
private let packetOutboundActor: PacketOutboundActor
|
||||
private let arpServer: ArpServer
|
||||
private let superServiceProxy: SDLSuperServiceProxy
|
||||
private let flowTracer: SDLFlowTracer
|
||||
|
||||
private var networkAddress: SDLConfiguration.NetworkAddress
|
||||
private var identityId: UInt32
|
||||
private var dataCipher: CCDataCipher?
|
||||
|
||||
init(provider: NEPacketTunnelProvider,
|
||||
config: SDLConfiguration,
|
||||
dataCipher: CCDataCipher?,
|
||||
policyService: PolicyService,
|
||||
packetOutboundActor: PacketOutboundActor,
|
||||
arpServer: ArpServer,
|
||||
superServiceProxy: SDLSuperServiceProxy,
|
||||
flowTracer: SDLFlowTracer) {
|
||||
self.provider = provider
|
||||
self.networkAddress = config.networkAddress
|
||||
self.identityId = config.identityId
|
||||
self.dataCipher = dataCipher
|
||||
self.policyService = policyService
|
||||
self.packetOutboundActor = packetOutboundActor
|
||||
self.arpServer = arpServer
|
||||
self.superServiceProxy = superServiceProxy
|
||||
self.flowTracer = flowTracer
|
||||
}
|
||||
|
||||
func updateRuntime(config: SDLConfiguration, dataCipher: CCDataCipher?) {
|
||||
self.networkAddress = config.networkAddress
|
||||
self.identityId = config.identityId
|
||||
self.dataCipher = dataCipher
|
||||
}
|
||||
|
||||
func handleData(_ data: SDLData) async {
|
||||
let processor = SDLHoleDataProcessor(
|
||||
networkAddress: self.networkAddress,
|
||||
dataCipher: self.dataCipher,
|
||||
policyService: self.policyService
|
||||
)
|
||||
|
||||
guard let plan = try? await processor.makeProcessingPlan(data: data) else {
|
||||
return
|
||||
}
|
||||
|
||||
self.flowTracer.inc(num: plan.inboundBytes, type: .inbound)
|
||||
|
||||
switch plan.action {
|
||||
case .sendARPReply(let dstMac, let responseData):
|
||||
SDLLogger.log("[PacketInboundActor] get arp request packet")
|
||||
await self.packetOutboundActor.routeLayerPacket(dstMac: dstMac, type: .arp, data: responseData)
|
||||
case .appendARP(let ip, let mac):
|
||||
SDLLogger.log("[PacketInboundActor] get arp response packet")
|
||||
await self.arpServer.append(ip: ip, mac: mac)
|
||||
case .writeToTun(let packetData, let identityID):
|
||||
let packet = NEPacket(data: packetData, protocolFamily: 2)
|
||||
self.provider.packetFlow.writePacketObjects([packet])
|
||||
SDLLogger.log("[PacketInboundActor] hole identity: \(identityID), allow, data count: \(packetData.count)", for: .trace)
|
||||
case .requestPolicy(let srcIdentityID):
|
||||
SDLLogger.log("[PacketInboundActor] not found identity: \(srcIdentityID) ruleMap", for: .debug)
|
||||
if let queryData = await self.policyService.identifyStore.makePolicyRequest(srcIdentityId: srcIdentityID, dstIdentityId: self.identityId) {
|
||||
await self.superServiceProxy.send(type: .policyRequest, data: queryData)
|
||||
}
|
||||
case .none:
|
||||
()
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -240,12 +240,28 @@ actor SDLUDPHoleService {
|
||||
}
|
||||
|
||||
actor SDLUDPHoleServiceProxy {
|
||||
typealias ControlEventHandler = @Sendable (SDLUDPHoleService.Event) async -> Void
|
||||
|
||||
private var udpHoleService: SDLUDPHoleService?
|
||||
private var packetInboundActor: PacketInboundActor?
|
||||
private var onControlEvent: ControlEventHandler?
|
||||
private var generation: UInt64 = 0
|
||||
|
||||
func replace(_ udpHoleService: SDLUDPHoleService?) async {
|
||||
self.generation &+= 1
|
||||
func bindInbound(_ packetInboundActor: PacketInboundActor, onControlEvent: @escaping ControlEventHandler) {
|
||||
self.packetInboundActor = packetInboundActor
|
||||
self.onControlEvent = onControlEvent
|
||||
}
|
||||
|
||||
func makeEventHandler() -> SDLUDPHoleService.EventHandler {
|
||||
self.generation &+= 1
|
||||
let generation = self.generation
|
||||
|
||||
return { [weak self] event in
|
||||
await self?.handleEvent(event, generation: generation)
|
||||
}
|
||||
}
|
||||
|
||||
func replace(_ udpHoleService: SDLUDPHoleService?) async {
|
||||
let oldUDPHoleService = self.udpHoleService
|
||||
self.udpHoleService = udpHoleService
|
||||
|
||||
@ -266,4 +282,22 @@ actor SDLUDPHoleServiceProxy {
|
||||
func send(type: SDLPacketType, data: Data, remoteAddress: SocketAddress) async {
|
||||
await self.udpHoleService?.send(type: type, data: data, remoteAddress: remoteAddress)
|
||||
}
|
||||
|
||||
private func handleEvent(_ event: SDLUDPHoleService.Event, generation: UInt64) async {
|
||||
guard generation == self.generation else {
|
||||
return
|
||||
}
|
||||
|
||||
switch event {
|
||||
case .packet(_, let message, _):
|
||||
switch message.inboundMessage {
|
||||
case .data(let data):
|
||||
await self.packetInboundActor?.handleData(data)
|
||||
case .control:
|
||||
await self.onControlEvent?(event)
|
||||
}
|
||||
case .ready, .natType, .closed:
|
||||
await self.onControlEvent?(event)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user