From 2de01a4c2942f2f86632cb8f09ff5196720ce7af Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=96=E7=95=8C?= Date: Wed, 8 Jul 2026 17:14:49 +0800 Subject: [PATCH] Reject flows on selector range exhaustion --- flow_dispatch.go | 36 ++++++++++++++++++++++++++++-------- 1 file changed, 28 insertions(+), 8 deletions(-) diff --git a/flow_dispatch.go b/flow_dispatch.go index 152a5bb..8911bf1 100644 --- a/flow_dispatch.go +++ b/flow_dispatch.go @@ -113,6 +113,7 @@ type ForwardDispatcher struct { activeNATs []*portNAT writebackBatch [][]byte returnPath forwardReturn + exhaustedLogAt int64 segmentBuffers [][]byte segmentSizes []int @@ -247,14 +248,23 @@ func (d *ForwardDispatcher) judgeAndInstall(key flowKey, packet *forwardPacket, switch verdict.Action { case ActionFlow: if verdict.Port != nil { - flow, created := d.createFlow(packet, verdict) - if created { + flow, result := d.createFlow(packet, verdict) + if result == createFlowOK { entry := &flowEntry{action: ActionFlow, flow: flow, idle: d.flowIdle(flow)} entry.deadline = now + int64(entry.idle) d.insertEntry(key, entry, now) d.forwardToPort(flow, packet, raw) return true } + if result == createFlowExhausted { + if now-d.exhaustedLogAt >= int64(exhaustedLogInterval) { + d.exhaustedLogAt = now + d.logger.Warn("port selector range exhausted, rejecting flow to ", packet.destination) + } + d.installSimple(key, ActionReject, packet.protocol, now) + d.stageReject(packet) + return true + } } d.installSimple(key, ActionAccept, packet.protocol, now) return false @@ -302,7 +312,17 @@ func (d *ForwardDispatcher) flowIdle(flow *forwardFlow) time.Duration { return d.idleTimeout(flow.protocol, established) } -func (d *ForwardDispatcher) createFlow(packet *forwardPacket, verdict FlowVerdict) (*forwardFlow, bool) { +type createFlowResult uint8 + +const ( + createFlowOK createFlowResult = iota + createFlowUnsupported + createFlowExhausted +) + +const exhaustedLogInterval = 5 * time.Second + +func (d *ForwardDispatcher) createFlow(packet *forwardPacket, verdict FlowVerdict) (*forwardFlow, createFlowResult) { var portAddress netip.Addr inet4Address, inet6Address := verdict.Port.PortAddresses() if packet.ipVersion == 6 { @@ -311,11 +331,11 @@ func (d *ForwardDispatcher) createFlow(packet *forwardPacket, verdict FlowVerdic portAddress = inet4Address } if !portAddress.IsValid() { - return nil, false + return nil, createFlowUnsupported } effectiveMTU := verdict.Port.PortMTU() if packet.ipVersion == 6 && effectiveMTU != 0 && effectiveMTU < header.IPv6MinimumMTU { - return nil, false + return nil, createFlowUnsupported } isICMP := isICMPProtocol(packet.protocol) clientDestinationAddress := packet.destination.Addr() @@ -330,11 +350,11 @@ func (d *ForwardDispatcher) createFlow(packet *forwardPacket, verdict FlowVerdic } nat := d.natFor(verdict.Port) if nat == nil { - return nil, false + return nil, createFlowUnsupported } selector, reverseKey, allocated := nat.allocateSelector(packet.protocol, portAddress, serverAddress, serverPort, packet.source.Port()) if !allocated { - return nil, false + return nil, createFlowExhausted } var udpTimeout time.Duration if packet.protocol == uint8(header.UDPProtocolNumber) { @@ -385,7 +405,7 @@ func (d *ForwardDispatcher) createFlow(packet *forwardPacket, verdict FlowVerdic } } nat.insert(reverseKey, flow) - return flow, true + return flow, createFlowOK } func (d *ForwardDispatcher) natFor(port Port) *portNAT {