Xray-core v1.260327.1-0.20260908222543-52a412d9e2f5 с изменениями Spanghew. Лицензия — Mozilla Public License 2.0 (полный текст — в разделе Xray-core). Исходный код неизменённого ядра: https://github.com/XTLS/Xray-core/tree/52a412d9e2f5 Ниже — изменённые файлы целиком и сами изменения (diff). Изменения, чьё имя начинается с «apple/», есть только в сборках для iPhone и Mac. ============================================================================== proxy/freedom/freedom.go (изменённый файл целиком) package freedom import ( "context" "crypto/rand" "io" "strings" "sync/atomic" "time" "github.com/pires/go-proxyproto" "github.com/xtls/xray-core/common" "github.com/xtls/xray-core/common/buf" "github.com/xtls/xray-core/common/crypto" "github.com/xtls/xray-core/common/dice" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/geodata" "github.com/xtls/xray-core/common/net" "github.com/xtls/xray-core/common/platform" "github.com/xtls/xray-core/common/retry" "github.com/xtls/xray-core/common/session" "github.com/xtls/xray-core/common/signal" "github.com/xtls/xray-core/common/task" "github.com/xtls/xray-core/common/utils" "github.com/xtls/xray-core/core" "github.com/xtls/xray-core/features/outbound" "github.com/xtls/xray-core/features/policy" "github.com/xtls/xray-core/features/routing" routing_session "github.com/xtls/xray-core/features/routing/session" "github.com/xtls/xray-core/features/stats" "github.com/xtls/xray-core/proxy" "github.com/xtls/xray-core/transport" "github.com/xtls/xray-core/transport/internet" "github.com/xtls/xray-core/transport/internet/stat" ) var ( useSplice atomic.Bool allNetworks [8]bool defaultBlockPrivateRule *FinalRule defaultBlockAllRule *FinalRule ) func reloadEnvSettings() error { const defaultFlagValue = "NOT_DEFINED_AT_ALL" value := platform.NewEnvFlag(platform.UseFreedomSplice).GetValue(func() string { return defaultFlagValue }) enabled := false switch value { case defaultFlagValue, "auto", "enable": enabled = true } useSplice.Store(enabled) return nil } func init() { common.Must(common.RegisterConfig((*Config)(nil), func(ctx context.Context, config interface{}) (interface{}, error) { h := new(Handler) h.instance = core.FromContext(ctx) if streamSettings, ok := session.StreamSettingsFromContext(ctx).(*internet.MemoryStreamConfig); ok && streamSettings.SocketSettings != nil { h.resolveStrategy = streamSettings.SocketSettings.DomainStrategy h.usesDialerProxy = len(streamSettings.SocketSettings.DialerProxy) > 0 } if err := core.RequireFeatures(ctx, func(pm policy.Manager) error { return h.Init(config.(*Config), pm) }); err != nil { return nil, err } return h, nil })) platform.RegisterEnvReload(reloadEnvSettings) for i := range allNetworks { allNetworks[i] = true } defaultBlockPrivateRule = &FinalRule{ action: RuleAction_Block, network: allNetworks, ip: geodata.GetPrivateIPMatcher(), } defaultBlockAllRule = &FinalRule{ action: RuleAction_Block, network: allNetworks, } } type FinalRule struct { action RuleAction network [8]bool port net.MemoryPortList ip geodata.IPMatcher blockDelay *Range } // Handler handles Freedom connections. type Handler struct { policyManager policy.Manager config *Config finalRules []*FinalRule resolveStrategy internet.DomainStrategy usesDialerProxy bool instance *core.Instance } // Spanghew: the route of a UDP connection is chosen by its first packet, and freedom then sent every // later packet of that socket to whatever address it carried. A socket whose first packet went to an // address the routing sends "direct" (a rule by address, port or sniffed protocol) carried its packets // to all other public addresses past the VPN as well: a STUN request of the same socket showed a // foreign server the real address of the device. Now, before a packet goes to an address other than // the one the connection was routed by, freedom asks the router where a connection to that address // would go (same inbound, source and sniffed protocol; no domain: it belonged to the first packet) and // drops the packet unless the answer is this outbound. type udpRouteGuard struct { router routing.Router ohm outbound.Manager inbound *session.Inbound content *session.Content tag string first net.Destination checked map[net.Destination]bool } // udpRouteGuardCache: destinations a connection remembers the answer for (a DHT socket talks to // thousands); full - forgotten and asked again. const udpRouteGuardCache = 256 func (h *Handler) newUDPRouteGuard(ctx context.Context, ob *session.Outbound) *udpRouteGuard { if h.instance == nil || h.usesDialerProxy || h.config.DestinationOverride != nil { // Not the final outbound, or the address is the one of the config: packets do not choose it. return nil } router, _ := h.instance.GetFeature(routing.RouterType()).(routing.Router) ohm, _ := h.instance.GetFeature(outbound.ManagerType()).(outbound.Manager) if router == nil || ohm == nil { return nil } guard := &udpRouteGuard{ router: router, ohm: ohm, inbound: session.InboundFromContext(ctx), tag: ob.Tag, first: ob.OriginalTarget, checked: make(map[net.Destination]bool), } if !guard.first.IsValid() { guard.first = ob.Target } if content := session.ContentFromContext(ctx); content != nil { guard.content = &session.Content{Protocol: content.Protocol, Attributes: content.Attributes} } return guard } // allows: requested - the address as the packet carried it, dest - the address it is about to be sent // to (the same, or the IP a domain was resolved to). func (g *udpRouteGuard) allows(requested, dest net.Destination) bool { if g == nil || requested == g.first || dest == g.first { return true } if allowed, found := g.checked[dest]; found { return allowed } var tag string route, err := g.router.PickRoute(&routing_session.Context{ Inbound: g.inbound, Outbound: &session.Outbound{OriginalTarget: dest, Target: dest}, Content: g.content, }) if err == nil { tag = route.GetOutboundTag() } else if handler := g.ohm.GetDefaultHandler(); handler != nil { tag = handler.Tag() } allowed := tag == g.tag if !allowed { errors.LogInfo(context.Background(), "UDP to ", dest, " is routed to [", tag, "], not to [", g.tag, "] the connection to ", g.first, " took: dropped") } if len(g.checked) >= udpRouteGuardCache { clear(g.checked) } g.checked[dest] = allowed return allowed } func buildFinalRule(config *FinalRuleConfig) (*FinalRule, error) { rule := &FinalRule{ action: config.GetAction(), blockDelay: config.GetBlockDelay(), } if len(config.Networks) == 0 { rule.network = allNetworks } else { for _, network := range config.Networks { rule.network[int(network)] = true } } if config.PortList != nil { rule.port = net.PortListFromProto(config.PortList) } if len(config.Ip) > 0 { matcher, err := geodata.IPReg.BuildIPMatcher(config.Ip) if err != nil { return nil, err } rule.ip = matcher } return rule, nil } func (r *FinalRule) matchNetwork(network net.Network) bool { return r.network[int(network)] } func (r *FinalRule) matchPort(port net.Port) bool { if len(r.port) == 0 { return true } return r.port.Contains(port) } func (r *FinalRule) matchIP(addr net.Address) bool { if r.ip == nil { return true } return addr != nil && addr.Family().IsIP() && r.ip.Match(addr.IP()) } func (r *FinalRule) Apply(network net.Network, address net.Address, port net.Port) bool { if !r.matchNetwork(network) { return false } if !r.matchPort(port) { return false } return r.matchIP(address) } func getDefaultFinalRule(inbound *session.Inbound) *FinalRule { if inbound == nil { return nil } switch inbound.Name { case "vless-reverse": return defaultBlockAllRule case "vless", "vmess", "trojan", "hysteria", "wireguard": return defaultBlockPrivateRule default: if strings.HasPrefix(inbound.Name, "shadowsocks") { return defaultBlockPrivateRule } } return nil } func (h *Handler) matchFinalRule(network net.Network, address net.Address, port net.Port, defaultRule *FinalRule) *FinalRule { for _, rule := range h.finalRules { if rule.Apply(network, address, port) { return rule } } if defaultRule != nil && defaultRule.Apply(network, address, port) { return defaultRule } return nil } // Init initializes the Handler with necessary parameters. func (h *Handler) Init(config *Config, pm policy.Manager) error { h.config = config h.policyManager = pm if h.usesDialerProxy { // freedom is not the final outbound, final rules do not apply if len(config.FinalRules) > 0 { errors.LogWarning(context.Background(), `The "finalRules" setting is ignored when "sockopt.dialerProxy" is set, since freedom is not the final outbound.`) } return nil } h.finalRules = make([]*FinalRule, 0, len(config.FinalRules)) for _, rc := range config.FinalRules { rule, err := buildFinalRule(rc) if err != nil { return errors.New("failed to build final rule").Base(err) } h.finalRules = append(h.finalRules, rule) } return nil } func (h *Handler) policy() policy.Session { p := h.policyManager.ForLevel(h.config.UserLevel) return p } func (h *Handler) blockDelay(rule *FinalRule) time.Duration { min := uint64(30) max := uint64(90) if rule.blockDelay != nil { min = rule.blockDelay.Min max = rule.blockDelay.Max } span := max - min if max < min { span = min - max } return time.Duration(min+uint64(dice.Roll(int(span+1)))) * time.Second } func (h *Handler) blackhole(ctx context.Context, input buf.Reader, output buf.Writer, rule *FinalRule, dest *net.Destination) error { delay := h.blockDelay(rule) errors.LogInfo(ctx, "blocked target: ", *dest, ", blackholing connection for ", delay) timer := time.AfterFunc(delay, func() { common.Interrupt(input) common.Interrupt(output) errors.LogInfo(ctx, "closed blackholed connection to blocked target: ", *dest) }) defer timer.Stop() defer common.Close(output) _ = buf.Copy(input, buf.Discard) return nil } func isValidAddress(addr *net.IPOrDomain) bool { if addr == nil { return false } a := addr.AsAddress() return a != net.AnyIP && a != net.AnyIPv6 } // Process implements proxy.Outbound. func (h *Handler) Process(ctx context.Context, link *transport.Link, dialer internet.Dialer) error { outbounds := session.OutboundsFromContext(ctx) ob := outbounds[len(outbounds)-1] if !ob.Target.IsValid() { return errors.New("target not specified.") } ob.Name = "freedom" ob.CanSpliceCopy = 1 inbound := session.InboundFromContext(ctx) var defaultRule *FinalRule if !h.usesDialerProxy { // freedom is not the final outbound, final rules do not apply (and the domain is not resolved) defaultRule = getDefaultFinalRule(inbound) } destination := ob.Target origTargetAddr := ob.OriginalTarget.Address if origTargetAddr == nil { origTargetAddr = ob.Target.Address } dialer.SetOutboundGateway(ctx, ob) outGateway := ob.Gateway UDPOverride := net.UDPDestination(nil, 0) if h.config.DestinationOverride != nil { server := h.config.DestinationOverride.Server if isValidAddress(server.Address) { destination.Address = server.Address.AsAddress() UDPOverride.Address = destination.Address } if server.Port != 0 { destination.Port = net.Port(server.Port) UDPOverride.Port = destination.Port } } input := link.Reader output := link.Writer var conn stat.Connection var blockedDest *net.Destination var blockedRule *FinalRule err := retry.ExponentialBackoff(5, 100).On(func() error { if destination.Address.Family().IsDomain() { if defaultRule != nil || len(h.finalRules) > 0 { if strategy := h.resolveStrategy; strategy.HasStrategy() { ips, err := internet.LookupForIP(destination.Address.Domain(), strategy, outGateway) if err != nil { // non-force may still dial with system DNS errors.LogInfoInner(ctx, err, "failed to get IP address for domain ", destination.Address.Domain()) if strategy.ForceIP() { return err // retry } } for _, ip := range ips { if addr := net.IPAddress(ip); addr != nil { if rule := h.matchFinalRule(destination.Network, addr, destination.Port, defaultRule); rule != nil && rule.action == RuleAction_Block { blockedDest = &destination blockedDest.Address = addr blockedRule = rule return nil } } } } else { addrs, err := net.DefaultResolver.LookupIPAddr(ctx, destination.Address.Domain()) if err != nil { // dialer may retry DNS errors.LogInfoInner(ctx, err, "failed to get IP address for domain ", destination.Address.Domain()) } for _, addr := range addrs { if ipAddr := net.IPAddress(addr.IP); ipAddr != nil { if rule := h.matchFinalRule(destination.Network, ipAddr, destination.Port, defaultRule); rule != nil && rule.action == RuleAction_Block { blockedDest = &destination blockedDest.Address = ipAddr blockedRule = rule return nil } } } } } } else { if rule := h.matchFinalRule(destination.Network, destination.Address, destination.Port, defaultRule); rule != nil && rule.action == RuleAction_Block { blockedDest = &destination blockedRule = rule return nil } } rawConn, err := dialer.Dial(ctx, destination) if err != nil { return err } conn = rawConn return nil }) if err != nil { return errors.New("failed to open connection to ", destination).Base(err) } if blockedDest != nil { return h.blackhole(ctx, input, output, blockedRule, blockedDest) } if destination.Address.Family().IsDomain() && (defaultRule != nil || len(h.finalRules) > 0) { // pre-check may fail or dialer may select another IP remoteDest := net.DestinationFromAddr(conn.RemoteAddr()) if rule := h.matchFinalRule(remoteDest.Network, remoteDest.Address, remoteDest.Port, defaultRule); rule != nil && rule.action == RuleAction_Block { conn.Close() return h.blackhole(ctx, input, output, rule, &remoteDest) } } if h.config.ProxyProtocol > 0 && h.config.ProxyProtocol <= 2 { version := byte(h.config.ProxyProtocol) srcAddr := inbound.Source.RawNetAddr() dstAddr := conn.RemoteAddr() header := proxyproto.HeaderProxyFromAddrs(version, srcAddr, dstAddr) if _, err = header.WriteTo(conn); err != nil { conn.Close() return errors.New("failed to set PROXY protocol v", version).Base(err) } } defer conn.Close() errors.LogInfo(ctx, "connection opened to ", destination, ", local endpoint ", conn.LocalAddr(), ", remote endpoint ", conn.RemoteAddr()) var newCtx context.Context var newCancel context.CancelFunc if session.TimeoutOnlyFromContext(ctx) { newCtx, newCancel = context.WithCancel(context.Background()) } plcy := h.policy() ctx, cancel := context.WithCancel(ctx) timer := signal.CancelAfterInactivity(ctx, func() { cancel() if newCancel != nil { newCancel() } }, plcy.Timeouts.ConnectionIdle) requestDone := func() error { defer timer.SetTimeout(plcy.Timeouts.DownlinkOnly) var writer buf.Writer if destination.Network == net.Network_TCP { if h.config.Fragment != nil { errors.LogDebug(ctx, "FRAGMENT", h.config.Fragment.PacketsFrom, h.config.Fragment.PacketsTo, h.config.Fragment.LengthMin, h.config.Fragment.LengthMax, h.config.Fragment.IntervalMin, h.config.Fragment.IntervalMax, h.config.Fragment.MaxSplitMin, h.config.Fragment.MaxSplitMax) writer = buf.NewWriter(&FragmentWriter{ fragment: h.config.Fragment, writer: conn, }) } else { writer = buf.NewWriter(conn) } } else { writer = NewPacketWriter(conn, h, defaultRule, UDPOverride, destination, outGateway) if packetWriter, ok := writer.(*PacketWriter); ok { packetWriter.RouteGuard = h.newUDPRouteGuard(ctx, ob) } if h.config.Noises != nil { errors.LogDebug(ctx, "NOISE", h.config.Noises) writer = &NoisePacketWriter{ Writer: writer, noises: h.config.Noises, firstWrite: true, UDPOverride: UDPOverride, remoteAddr: net.DestinationFromAddr(conn.RemoteAddr()).Address, } } } if err := buf.Copy(input, writer, buf.UpdateActivity(timer)); err != nil { return errors.New("failed to process request").Base(err) } return nil } responseDone := func() error { defer timer.SetTimeout(plcy.Timeouts.UplinkOnly) if destination.Network == net.Network_TCP && useSplice.Load() && proxy.IsRAWTransportWithoutSecurity(conn) { // it would be tls conn in special use case of MITM, we need to let link handle traffic var writeConn net.Conn var inTimer *signal.ActivityTimer if inbound := session.InboundFromContext(ctx); inbound != nil && inbound.Conn != nil { writeConn = inbound.Conn inTimer = inbound.Timer } return proxy.CopyRawConnIfExist(ctx, conn, writeConn, link.Writer, timer, inTimer) } var reader buf.Reader if destination.Network == net.Network_TCP { reader = buf.NewReader(conn) } else { reader = NewPacketReader(conn, h, defaultRule, UDPOverride, destination) } if err := buf.Copy(reader, output, buf.UpdateActivity(timer)); err != nil { return errors.New("failed to process response").Base(err) } return nil } if newCtx != nil { ctx = newCtx } if err := task.Run(ctx, requestDone, task.OnSuccess(responseDone, task.Close(output))); err != nil { return errors.New("connection ends").Base(err) } return nil } func NewPacketReader(conn net.Conn, h *Handler, defaultRule *FinalRule, UDPOverride net.Destination, DialDest net.Destination) buf.Reader { iConn := conn statConn, ok := iConn.(*stat.CounterConnection) if ok { iConn = statConn.Connection } var counter stats.Counter if statConn != nil { counter = statConn.ReadCounter } if c, ok := iConn.(*internet.PacketConnWrapper); ok { isOverridden := false if UDPOverride.Address != nil || UDPOverride.Port != 0 { isOverridden = true } return &PacketReader{ PacketConnWrapper: c, Counter: counter, Handler: h, DefaultRule: defaultRule, IsOverridden: isOverridden, InitUnchangedAddr: DialDest.Address, InitChangedAddr: net.DestinationFromAddr(conn.RemoteAddr()).Address, } } return &buf.PacketReader{Reader: conn} } type PacketReader struct { *internet.PacketConnWrapper stats.Counter Handler *Handler DefaultRule *FinalRule IsOverridden bool InitUnchangedAddr net.Address InitChangedAddr net.Address } func (r *PacketReader) ReadMultiBuffer() (buf.MultiBuffer, error) { b := buf.New() b.Resize(0, buf.Size) for { n, d, err := r.PacketConnWrapper.ReadFrom(b.Bytes()) if err != nil { b.Release() return nil, err } udpAddr := d.(*net.UDPAddr) sourceAddr := net.IPAddress(udpAddr.IP) if rule := r.Handler.matchFinalRule(net.Network_UDP, sourceAddr, net.Port(udpAddr.Port), r.DefaultRule); rule != nil && rule.action == RuleAction_Block { continue } b.Resize(0, int32(n)) // if udp dest addr is changed, we are unable to get the correct src addr // so we don't attach src info to udp packet, break cone behavior, assuming the dial dest is the expected scr addr if !r.IsOverridden { if r.InitChangedAddr == sourceAddr { sourceAddr = r.InitUnchangedAddr } b.UDP = &net.Destination{ Address: sourceAddr, Port: net.Port(udpAddr.Port), Network: net.Network_UDP, } } if r.Counter != nil { r.Counter.Add(int64(n)) } return buf.MultiBuffer{b}, nil } } // DialDest means the dial target used in the dialer when creating conn func NewPacketWriter(conn net.Conn, h *Handler, defaultRule *FinalRule, UDPOverride net.Destination, DialDest net.Destination, outGateway net.Address) buf.Writer { iConn := conn statConn, ok := iConn.(*stat.CounterConnection) if ok { iConn = statConn.Connection } var counter stats.Counter if statConn != nil { counter = statConn.WriteCounter } if c, ok := iConn.(*internet.PacketConnWrapper); ok { // If DialDest is a domain, it will be resolved in dialer // check this behavior and add it to map resolvedUDPAddr := utils.NewTypedSyncMap[string, net.Address]() if DialDest.Address.Family().IsDomain() { resolvedUDPAddr.Store(DialDest.Address.Domain(), net.DestinationFromAddr(conn.RemoteAddr()).Address) } return &PacketWriter{ PacketConnWrapper: c, Counter: counter, Handler: h, DefaultRule: defaultRule, UDPOverride: UDPOverride, ResolvedUDPAddr: resolvedUDPAddr, OutGateway: outGateway, } } return &buf.SequentialWriter{Writer: conn} } type PacketWriter struct { *internet.PacketConnWrapper stats.Counter *Handler DefaultRule *FinalRule UDPOverride net.Destination // Dest of udp packets might be a domain, we will resolve them to IP // But resolver will return a random one if the domain has many IPs // Resulting in these packets being sent to many different IPs randomly // So, cache and keep the resolve result ResolvedUDPAddr *utils.TypedSyncMap[string, net.Address] OutGateway net.Address // Spanghew: see udpRouteGuard; nil - no check. RouteGuard *udpRouteGuard } func (w *PacketWriter) WriteMultiBuffer(mb buf.MultiBuffer) error { for { mb2, b := buf.SplitFirst(mb) mb = mb2 if b == nil { break } var n int var err error if b.UDP != nil { requested := *b.UDP if w.UDPOverride.Address != nil { b.UDP.Address = w.UDPOverride.Address } if w.UDPOverride.Port != 0 { b.UDP.Port = w.UDPOverride.Port } if b.UDP.Address.Family().IsDomain() { if ip, ok := w.ResolvedUDPAddr.Load(b.UDP.Address.Domain()); ok { b.UDP.Address = ip } else { shouldUseSystemResolver := true if strategy := w.Handler.resolveStrategy; strategy.HasStrategy() { ips, err := internet.LookupForIP(b.UDP.Address.Domain(), strategy, w.OutGateway) if err != nil { // drop packet if resolve failed when forceIP if strategy.ForceIP() { b.Release() continue } } else { ip = net.IPAddress(ips[dice.Roll(len(ips))]) shouldUseSystemResolver = false } } if shouldUseSystemResolver { udpAddr, err := net.ResolveUDPAddr("udp", b.UDP.NetAddr()) if err != nil { b.Release() continue } else { ip = net.IPAddress(udpAddr.IP) } } if ip != nil { b.UDP.Address, _ = w.ResolvedUDPAddr.LoadOrStore(b.UDP.Address.Domain(), ip) } } } if rule := w.matchFinalRule(net.Network_UDP, b.UDP.Address, b.UDP.Port, w.DefaultRule); rule != nil && rule.action == RuleAction_Block { b.Release() continue } if !w.RouteGuard.allows(requested, *b.UDP) { b.Release() continue } destAddr := b.UDP.RawNetAddr() if destAddr == nil { b.Release() continue } n, err = w.PacketConnWrapper.WriteTo(b.Bytes(), destAddr) } else { n, err = w.PacketConnWrapper.Write(b.Bytes()) } b.Release() if err != nil { buf.ReleaseMulti(mb) return err } if w.Counter != nil { w.Counter.Add(int64(n)) } } return nil } type NoisePacketWriter struct { buf.Writer noises []*Noise firstWrite bool UDPOverride net.Destination remoteAddr net.Address } // MultiBuffer writer with Noise before first packet func (w *NoisePacketWriter) WriteMultiBuffer(mb buf.MultiBuffer) error { if w.firstWrite { w.firstWrite = false // Do not send Noise for dns requests(just to be safe) if w.UDPOverride.Port == 53 { return w.Writer.WriteMultiBuffer(mb) } var noise []byte var err error if w.remoteAddr.Family().IsDomain() { panic("impossible, remoteAddr is always IP") } for _, n := range w.noises { switch n.ApplyTo { case "ipv4": if w.remoteAddr.Family().IsIPv6() { continue } case "ipv6": if w.remoteAddr.Family().IsIPv4() { continue } case "ip": default: panic("unreachable, applyTo is ip/ipv4/ipv6") } // User input string or base64 encoded string or hex string if n.Packet != nil { noise = n.Packet } else { // Random noise noise, err = GenerateRandomBytes(crypto.RandBetween(int64(n.LengthMin), int64(n.LengthMax))) } if err != nil { return err } err = w.Writer.WriteMultiBuffer(buf.MultiBuffer{buf.FromBytes(noise)}) if err != nil { return err } if n.DelayMin != 0 || n.DelayMax != 0 { time.Sleep(time.Duration(crypto.RandBetween(int64(n.DelayMin), int64(n.DelayMax))) * time.Millisecond) } } } return w.Writer.WriteMultiBuffer(mb) } type FragmentWriter struct { fragment *Fragment writer io.Writer count uint64 } func (f *FragmentWriter) Write(b []byte) (int, error) { f.count++ if f.fragment.PacketsFrom == 0 && f.fragment.PacketsTo == 1 { if f.count != 1 || len(b) <= 5 || b[0] != 22 { return f.writer.Write(b) } recordLen := 5 + ((int(b[3]) << 8) | int(b[4])) if len(b) < recordLen { // maybe already fragmented somehow return f.writer.Write(b) } data := b[5:recordLen] buff := make([]byte, 2048) var hello []byte maxSplit := crypto.RandBetween(int64(f.fragment.MaxSplitMin), int64(f.fragment.MaxSplitMax)) var splitNum int64 for from := 0; ; { to := from + int(crypto.RandBetween(int64(f.fragment.LengthMin), int64(f.fragment.LengthMax))) splitNum++ if to > len(data) || (maxSplit > 0 && splitNum >= maxSplit) { to = len(data) } l := to - from if 5+l > len(buff) { buff = make([]byte, 5+l) } copy(buff[:3], b) copy(buff[5:], data[from:to]) from = to buff[3] = byte(l >> 8) buff[4] = byte(l) if f.fragment.IntervalMax == 0 { // combine fragmented tlshello if interval is 0 hello = append(hello, buff[:5+l]...) } else { _, err := f.writer.Write(buff[:5+l]) time.Sleep(time.Duration(crypto.RandBetween(int64(f.fragment.IntervalMin), int64(f.fragment.IntervalMax))) * time.Millisecond) if err != nil { return 0, err } } if from == len(data) { if len(hello) > 0 { _, err := f.writer.Write(hello) if err != nil { return 0, err } } if len(b) > recordLen { n, err := f.writer.Write(b[recordLen:]) if err != nil { return recordLen + n, err } } return len(b), nil } } } if f.fragment.PacketsFrom != 0 && (f.count < f.fragment.PacketsFrom || f.count > f.fragment.PacketsTo) { return f.writer.Write(b) } maxSplit := crypto.RandBetween(int64(f.fragment.MaxSplitMin), int64(f.fragment.MaxSplitMax)) var splitNum int64 for from := 0; ; { to := from + int(crypto.RandBetween(int64(f.fragment.LengthMin), int64(f.fragment.LengthMax))) splitNum++ if to > len(b) || (maxSplit > 0 && splitNum >= maxSplit) { to = len(b) } n, err := f.writer.Write(b[from:to]) from += n if err != nil { return from, err } time.Sleep(time.Duration(crypto.RandBetween(int64(f.fragment.IntervalMin), int64(f.fragment.IntervalMax))) * time.Millisecond) if from >= len(b) { return from, nil } } } func GenerateRandomBytes(n int64) ([]byte, error) { b := make([]byte, n) _, err := rand.Read(b) // Note that err == nil only if we read len(b) bytes. if err != nil { return nil, err } return b, nil } ============================================================================== proxy/freedom/spanghew_udp_route_test.go (изменённый файл целиком) package freedom import ( "context" "testing" "time" "github.com/xtls/xray-core/app/router" "github.com/xtls/xray-core/common" "github.com/xtls/xray-core/common/buf" "github.com/xtls/xray-core/common/geodata" "github.com/xtls/xray-core/common/net" "github.com/xtls/xray-core/common/serial" "github.com/xtls/xray-core/common/session" "github.com/xtls/xray-core/features/outbound" "github.com/xtls/xray-core/transport" "github.com/xtls/xray-core/transport/internet" ) type spanghewHandler struct{ tag string } func (h *spanghewHandler) Start() error { return nil } func (h *spanghewHandler) Close() error { return nil } func (h *spanghewHandler) Tag() string { return h.tag } func (h *spanghewHandler) Dispatch(context.Context, *transport.Link) {} func (h *spanghewHandler) SenderSettings() *serial.TypedMessage { return nil } func (h *spanghewHandler) ProxySettings() *serial.TypedMessage { return nil } // The first outbound ("proxy") is the default one, as in a subscription config. type spanghewOutbounds struct { outbound.Manager outbound.HandlerSelector } func (spanghewOutbounds) GetDefaultHandler() outbound.Handler { return &spanghewHandler{tag: "proxy"} } func spanghewCIDR(ip string, prefix uint32) *geodata.IPRule { return &geodata.IPRule{Value: &geodata.IPRule_Custom{Custom: &geodata.CIDRRule{ Cidr: &geodata.CIDR{Ip: net.ParseAddress(ip).IP(), Prefix: prefix}, }}} } // Rules of the shape a subscription template has: torrents and "home" addresses - direct, some of // the home addresses - to the server all the same (the rule stands earlier), the rest - the default // outbound. func spanghewRouter(t *testing.T) *router.Router { t.Helper() r := new(router.Router) common.Must(r.Init(context.TODO(), &router.Config{ Rule: []*router.RoutingRule{ { TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, Protocol: []string{"bittorrent"}, }, { TargetTag: &router.RoutingRule_Tag{Tag: "proxy"}, Ip: []*geodata.IPRule{spanghewCIDR("203.0.113.128", 25)}, }, { TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, Ip: []*geodata.IPRule{spanghewCIDR("203.0.113.0", 24)}, }, { TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, Networks: []net.Network{net.Network_UDP}, PortList: &net.PortList{Range: []*net.PortRange{{From: 6881, To: 6889}}}, }, { TargetTag: &router.RoutingRule_Tag{Tag: "other-direct"}, Ip: []*geodata.IPRule{spanghewCIDR("192.168.0.0", 16)}, }, }, }, nil, &spanghewOutbounds{}, nil)) return r } func spanghewGuard(t *testing.T, first net.Destination, protocol string) *udpRouteGuard { t.Helper() return &udpRouteGuard{ router: spanghewRouter(t), ohm: &spanghewOutbounds{}, inbound: &session.Inbound{Tag: "tun-in", Source: net.UDPDestination(net.ParseAddress("172.19.0.1"), 50000)}, content: &session.Content{Protocol: protocol}, tag: "direct", first: first, checked: make(map[net.Destination]bool), } } func spanghewUDP(address string, port net.Port) net.Destination { return net.UDPDestination(net.ParseAddress(address), port) } // The reported leak: the first packet of a socket goes to an address routed "direct", the STUN // request of the same socket - to a foreign server. The second must not leave through "direct". func TestSpanghewUDPRouteGuard(t *testing.T) { first := spanghewUDP("203.0.113.8", 9) guard := spanghewGuard(t, first, "") for dest, allowed := range map[net.Destination]bool{ first: true, spanghewUDP("203.0.113.9", 3478): true, // another "home" address: the same route spanghewUDP("203.0.113.8", 3478): true, // another port of the first address spanghewUDP("162.159.207.0", 3478): false, // a foreign address: the default outbound spanghewUDP("8.8.8.8", 53): false, spanghewUDP("203.0.113.200", 3478): false, // "home", but the earlier rule sends it to the server spanghewUDP("192.168.1.1", 53): false, // direct too, but through another outbound spanghewUDP("162.159.207.0", 6881): true, // the rule by port spanghewUDP("2606:4700::1111", 443): false, } { for range 2 { // the second answer is the remembered one if got := guard.allows(dest, dest); got != allowed { t.Errorf("%v: allowed %v, want %v", dest, got, allowed) } } } } // The destination the connection was routed by is not asked about: its route may have been chosen // by a sniffed domain, which the packets do not carry. func TestSpanghewUDPRouteGuardFirstDestination(t *testing.T) { first := spanghewUDP("162.159.207.0", 443) guard := spanghewGuard(t, first, "quic") if !guard.allows(first, first) { t.Error("the first destination is dropped") } if !guard.allows(net.UDPDestination(net.DomainAddress("example.com"), 443), first) { t.Error("the first destination, resolved, is dropped") } if guard.allows(spanghewUDP("162.159.207.0", 3478), spanghewUDP("162.159.207.0", 3478)) { t.Error("another port of the first address went past the routing") } if guard.allows(spanghewUDP("162.159.207.1", 443), spanghewUDP("162.159.207.1", 443)) { t.Error("another address went past the routing") } } // A rule by sniffed protocol holds for the whole connection: torrents routed "direct" reach all // their peers. func TestSpanghewUDPRouteGuardKeepsProtocol(t *testing.T) { guard := spanghewGuard(t, spanghewUDP("162.159.207.0", 51413), "bittorrent") for _, dest := range []net.Destination{spanghewUDP("8.8.8.8", 6000), spanghewUDP("203.0.113.200", 6000)} { if !guard.allows(dest, dest) { t.Errorf("%v: a torrent peer is dropped", dest) } } } // Without a guard (not the final outbound, an address override) nothing changes. func TestSpanghewUDPRouteGuardAbsent(t *testing.T) { var guard *udpRouteGuard dest := spanghewUDP("8.8.8.8", 53) if !guard.allows(dest, dest) { t.Error("a nil guard drops") } handler := &Handler{config: &Config{}} if handler.newUDPRouteGuard(context.Background(), &session.Outbound{Target: dest}) != nil { t.Error("a guard without a core instance") } } func TestSpanghewUDPRouteGuardForgets(t *testing.T) { guard := spanghewGuard(t, spanghewUDP("203.0.113.8", 9), "") for port := 1; port <= 3*udpRouteGuardCache; port++ { dest := spanghewUDP("8.8.8.8", net.Port(port)) if guard.allows(dest, dest) { t.Fatalf("%v went past the routing", dest) } if len(guard.checked) > udpRouteGuardCache { t.Fatalf("remembered %d destinations", len(guard.checked)) } } } // Through the writer itself: of two packets of one connection only the one to the address routed // here reaches the network. func TestSpanghewPacketWriterDropsOtherRoute(t *testing.T) { listen := func() *net.UDPConn { conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.IP{127, 0, 0, 1}}) common.Must(err) t.Cleanup(func() { conn.Close() }) return conn } home, foreign, local := listen(), listen(), listen() homeDest := net.DestinationFromAddr(home.LocalAddr()) foreignDest := net.DestinationFromAddr(foreign.LocalAddr()) r := new(router.Router) common.Must(r.Init(context.TODO(), &router.Config{ Rule: []*router.RoutingRule{{ TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, PortList: &net.PortList{Range: []*net.PortRange{{From: uint32(homeDest.Port), To: uint32(homeDest.Port)}}}, }}, }, nil, &spanghewOutbounds{}, nil)) writer := &PacketWriter{ PacketConnWrapper: &internet.PacketConnWrapper{PacketConn: local, Dest: home.LocalAddr()}, Handler: &Handler{config: &Config{}}, UDPOverride: net.UDPDestination(nil, 0), RouteGuard: &udpRouteGuard{ router: r, ohm: &spanghewOutbounds{}, tag: "direct", first: homeDest, checked: make(map[net.Destination]bool), }, } packet := func(dest net.Destination, payload byte) *buf.Buffer { b := buf.New() b.WriteByte(payload) b.UDP = &dest return b } common.Must(writer.WriteMultiBuffer(buf.MultiBuffer{ packet(homeDest, 1), packet(foreignDest, 2), packet(homeDest, 3), })) read := func(conn *net.UDPConn) []byte { var got []byte data := make([]byte, 16) for { conn.SetReadDeadline(time.Now().Add(300 * time.Millisecond)) n, _, err := conn.ReadFrom(data) if err != nil { return got } got = append(got, data[:n]...) } } if got := read(home); len(got) != 2 || got[0] != 1 || got[1] != 3 { t.Errorf("the routed address received %v", got) } if got := read(foreign); len(got) != 0 { t.Errorf("the address of another route received %v", got) } } ============================================================================== infra/conf/shadowsocks.go (изменённый файл целиком) package conf import ( "strings" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/protocol" "github.com/xtls/xray-core/common/serial" "github.com/xtls/xray-core/common/task" "github.com/xtls/xray-core/proxy/shadowsocks" "google.golang.org/protobuf/proto" ) // Spanghew: Shadowsocks 2022 is cut out of this build. It was the only user of // github.com/sagernet/sing and sing-shadowsocks (GPLv3) through // proxy/shadowsocks_2022 and common/singbridge, and the client is closed source // (Android audit 2026-10-03, #5). Those two packages stay in the tree but nothing // imports them, so they are not linked. The methods get a clear error instead of // "unknown cipher method". func isShadowsocks2022(cipher string) bool { switch cipher { case "2022-blake3-aes-128-gcm", "2022-blake3-aes-256-gcm", "2022-blake3-chacha20-poly1305": return true } return false } var errShadowsocks2022 = errors.New("Shadowsocks 2022 is not included in this build") func cipherFromString(c string) shadowsocks.CipherType { switch strings.ToLower(c) { case "aes-128-gcm", "aead_aes_128_gcm": return shadowsocks.CipherType_AES_128_GCM case "aes-256-gcm", "aead_aes_256_gcm": return shadowsocks.CipherType_AES_256_GCM case "chacha20-poly1305", "aead_chacha20_poly1305", "chacha20-ietf-poly1305": return shadowsocks.CipherType_CHACHA20_POLY1305 case "xchacha20-poly1305", "aead_xchacha20_poly1305", "xchacha20-ietf-poly1305": return shadowsocks.CipherType_XCHACHA20_POLY1305 default: return shadowsocks.CipherType_UNKNOWN } } type ShadowsocksUserConfig struct { Cipher string `json:"method"` Password string `json:"password"` Level byte `json:"level"` Email string `json:"email"` Address *Address `json:"address"` Port uint16 `json:"port"` } type ShadowsocksServerConfig struct { Cipher string `json:"method"` Password string `json:"password"` Level byte `json:"level"` Email string `json:"email"` Users []*ShadowsocksUserConfig `json:"users"` Clients []*ShadowsocksUserConfig `json:"clients"` NetworkList *NetworkList `json:"network"` } func (v *ShadowsocksServerConfig) Build() (proto.Message, error) { errors.PrintNonRemovalDeprecatedFeatureWarning("Shadowsocks (with no Forward Secrecy, etc.)", "VLESS Encryption") if v.Clients != nil { v.Users = v.Clients } if isShadowsocks2022(v.Cipher) { return nil, errShadowsocks2022 } config := new(shadowsocks.ServerConfig) config.Network = v.NetworkList.Build() if v.Users != nil { if len(v.Users) > 0 { config.Users = make([]*protocol.User, len(v.Users)) processUser := func(idx int) error { user := v.Users[idx] account := &shadowsocks.Account{ Password: user.Password, CipherType: cipherFromString(user.Cipher), } if account.Password == "" { return errors.New("Shadowsocks password is not specified.") } if account.CipherType < shadowsocks.CipherType_AES_128_GCM || account.CipherType > shadowsocks.CipherType_XCHACHA20_POLY1305 { return errors.New("unsupported cipher method: ", user.Cipher) } config.Users[idx] = &protocol.User{ Email: user.Email, Level: uint32(user.Level), Account: serial.ToTypedMessage(account), } return nil } if err := task.ParallelForN(len(v.Users), processUser); err != nil { return nil, err } } } else { account := &shadowsocks.Account{ Password: v.Password, CipherType: cipherFromString(v.Cipher), } if account.Password == "" { return nil, errors.New("Shadowsocks password is not specified.") } if account.CipherType == shadowsocks.CipherType_UNKNOWN { return nil, errors.New("unknown cipher method: ", v.Cipher) } config.Users = append(config.Users, &protocol.User{ Email: v.Email, Level: uint32(v.Level), Account: serial.ToTypedMessage(account), }) } return config, nil } type ShadowsocksServerTarget struct { Address *Address `json:"address"` Port uint16 `json:"port"` Level byte `json:"level"` Email string `json:"email"` Cipher string `json:"method"` Password string `json:"password"` } type ShadowsocksClientConfig struct { Address *Address `json:"address"` Port uint16 `json:"port"` Level byte `json:"level"` Email string `json:"email"` Cipher string `json:"method"` Password string `json:"password"` Servers []*ShadowsocksServerTarget `json:"servers"` } func (v *ShadowsocksClientConfig) Build() (proto.Message, error) { errors.PrintNonRemovalDeprecatedFeatureWarning("Shadowsocks (with no Forward Secrecy, etc.)", "VLESS Encryption") if v.Address != nil { v.Servers = []*ShadowsocksServerTarget{ { Address: v.Address, Port: v.Port, Level: v.Level, Email: v.Email, Cipher: v.Cipher, Password: v.Password, }, } } if len(v.Servers) != 1 { return nil, errors.New(`Shadowsocks settings: "servers" should have one and only one member. Multiple endpoints in "servers" should use multiple Shadowsocks outbounds and routing balancer instead`) } config := new(shadowsocks.ClientConfig) for _, server := range v.Servers { if isShadowsocks2022(server.Cipher) { return nil, errShadowsocks2022 } if server.Address == nil { return nil, errors.New("Shadowsocks server address is not set.") } if server.Port == 0 { return nil, errors.New("Invalid Shadowsocks port.") } if server.Password == "" { return nil, errors.New("Shadowsocks password is not specified.") } account := &shadowsocks.Account{ Password: server.Password, } account.CipherType = cipherFromString(server.Cipher) if account.CipherType == shadowsocks.CipherType_UNKNOWN { return nil, errors.New("unknown cipher method: ", server.Cipher) } ss := &protocol.ServerEndpoint{ Address: server.Address.Build(), Port: uint32(server.Port), User: &protocol.User{ Level: uint32(server.Level), Email: server.Email, Account: serial.ToTypedMessage(account), }, } config.Server = ss break } return config, nil } ============================================================================== main/commands/all/api/inbound_user_add.go (изменённый файл целиком) package api import ( "context" "fmt" "github.com/xtls/xray-core/common/protocol" handlerService "github.com/xtls/xray-core/app/proxyman/command" cserial "github.com/xtls/xray-core/common/serial" "github.com/xtls/xray-core/core" "github.com/xtls/xray-core/infra/conf" "github.com/xtls/xray-core/infra/conf/serial" "github.com/xtls/xray-core/proxy/shadowsocks" "github.com/xtls/xray-core/proxy/trojan" vlessin "github.com/xtls/xray-core/proxy/vless/inbound" vmessin "github.com/xtls/xray-core/proxy/vmess/inbound" "github.com/xtls/xray-core/main/commands/base" ) var cmdAddInboundUsers = &base.Command{ CustomFlags: true, UsageLine: "{{.Exec}} api adu [--server=127.0.0.1:8080] [c2.json]...", Short: "Add users to inbounds", Long: ` Add users to inbounds. Arguments: -s, -server The API server address. Default 127.0.0.1:8080 -t, -timeout Timeout seconds to call API. Default 3 Example: {{.Exec}} {{.LongName}} --server=127.0.0.1:8080 c1.json c2.json `, Run: executeAddInboundUsers, } func executeAddInboundUsers(cmd *base.Command, args []string) { setSharedFlags(cmd) cmd.Flag.Parse(args) unnamedArgs := cmd.Flag.Args() inbs := extractInboundsConfig(unnamedArgs) conn, ctx, close := dialAPIServer() defer close() client := handlerService.NewHandlerServiceClient(conn) success := 0 for _, inb := range inbs { success += executeInboundUserAction(ctx, client, inb, addInboundUserAction) } fmt.Println("Added", success, "user(s) in total.") } func addInboundUserAction(ctx context.Context, client handlerService.HandlerServiceClient, tag string, user *protocol.User) error { fmt.Println("add user:", user.Email) _, err := client.AlterInbound(ctx, &handlerService.AlterInboundRequest{ Tag: tag, Operation: cserial.ToTypedMessage( &handlerService.AddUserOperation{ User: user, }, ), }) return err } func extractInboundUsers(inb *core.InboundHandlerConfig) []*protocol.User { if inb == nil { return nil } inst, err := inb.ProxySettings.GetInstance() if err != nil || inst == nil { fmt.Println("failed to get inbound instance:", err) return nil } switch ty := inst.(type) { case *vmessin.Config: return ty.User case *vlessin.Config: return ty.Users case *trojan.ServerConfig: return ty.Users case *shadowsocks.ServerConfig: return ty.Users default: fmt.Println("unsupported inbound type") } return nil } func extractInboundsConfig(unnamedArgs []string) []conf.InboundDetourConfig { ins := make([]conf.InboundDetourConfig, 0) for _, arg := range unnamedArgs { r, err := loadArg(arg) if err != nil { base.Fatalf("failed to load %s: %s", arg, err) } conf, err := serial.DecodeJSONConfig(r) if err != nil { base.Fatalf("failed to decode %s: %s", arg, err) } ins = append(ins, conf.InboundConfigs...) } return ins } func executeInboundUserAction(ctx context.Context, client handlerService.HandlerServiceClient, inb conf.InboundDetourConfig, action func(ctx context.Context, client handlerService.HandlerServiceClient, tag string, user *protocol.User) error) int { success := 0 tag := inb.Tag if len(tag) < 1 { return success } fmt.Println("processing inbound:", tag) built, err := inb.Build() if err != nil { fmt.Println("failed to build config:", err) return success } users := extractInboundUsers(built) if users == nil { return success } for _, user := range users { if len(user.Email) < 1 { continue } if err := action(ctx, client, inb.Tag, user); err == nil { fmt.Println("result: ok") success += 1 } else { fmt.Println(err) } } return success } ============================================================================== transport/internet/dialer.go (изменённый файл целиком) package internet import ( "context" "fmt" "strings" "sync/atomic" "github.com/xtls/xray-core/common" "github.com/xtls/xray-core/common/dice" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/net" "github.com/xtls/xray-core/common/net/cnc" "github.com/xtls/xray-core/common/session" "github.com/xtls/xray-core/features/dns" "github.com/xtls/xray-core/features/outbound" "github.com/xtls/xray-core/transport" "github.com/xtls/xray-core/transport/internet/stat" "github.com/xtls/xray-core/transport/pipe" ) // Dialer is the interface for dialing outbound connections. type Dialer interface { // Dial dials a system connection to the given destination. Dial(ctx context.Context, destination net.Destination) (stat.Connection, error) // DestIpAddress returns the ip of proxy server. It is useful in case of Android client, which prepare an IP before proxy connection is established DestIpAddress() net.IP // SetOutboundGateway set outbound gateway SetOutboundGateway(ctx context.Context, ob *session.Outbound) } // dialFunc is an interface to dial network connection to a specific destination. type dialFunc func(ctx context.Context, dest net.Destination, streamSettings *MemoryStreamConfig) (stat.Connection, error) var transportDialerCache = make(map[string]dialFunc) // RegisterTransportDialer registers a Dialer with given name. func RegisterTransportDialer(protocol string, dialer dialFunc) error { if _, found := transportDialerCache[protocol]; found { return errors.New(protocol, " dialer already registered").AtError() } transportDialerCache[protocol] = dialer return nil } // Dial dials a internet connection towards the given destination. func Dial(ctx context.Context, dest net.Destination, streamSettings *MemoryStreamConfig) (stat.Connection, error) { if dest.Network == net.Network_TCP { if streamSettings == nil { s, err := ToMemoryStreamConfig(nil) if err != nil { return nil, errors.New("failed to create default stream settings").Base(err) } streamSettings = s } protocol := streamSettings.ProtocolName dialer := transportDialerCache[protocol] if dialer == nil { return nil, errors.New(protocol, " dialer not registered").AtError() } return dialer(ctx, dest, streamSettings) } if dest.Network == net.Network_UDP { udpDialer := transportDialerCache["udp"] if udpDialer == nil { return nil, errors.New("UDP dialer not registered").AtError() } return udpDialer(ctx, dest, streamSettings) } return nil, errors.New("unknown network ", dest.Network) } // DestIpAddress returns the ip of proxy server. It is useful in case of Android client, which prepare an IP before proxy connection is established func DestIpAddress() net.IP { return effectiveSystemDialer.DestIpAddress() } // Spanghew: the DNS client and the outbound manager of the system dialer belong to the whole // process, and every core.New replaced them - two plain variables written while the connections of // a running core were reading them. The delay measurement of libXray creates its instance next to // the running VPN core (Android, Windows): for a moment the direct outbounds of that core asked the // resolver of the measurement (the names went past the tunnel) and found no outbound for // dialerProxy, and a reader could see the client of one instance with the manager of the other. // Now the pair is one value replaced atomically, and it can be held: while it is held, // InitSystemDialer of another instance changes nothing. type systemDialerState struct { dnsClient dns.Client obm outbound.Manager held bool } var systemDialer atomic.Pointer[systemDialerState] func LookupForIP(domain string, strategy DomainStrategy, localAddr net.Address) ([]net.IP, error) { var dnsClient dns.Client if state := systemDialer.Load(); state != nil { dnsClient = state.dnsClient } if dnsClient == nil { return nil, errors.New("DNS client not initialized").AtError() } ips, _, err := dnsClient.LookupIP(domain, dns.IPOption{ IPv4Enable: (localAddr == nil && strategy.PreferIP4()) || (localAddr != nil && localAddr.Family().IsIPv4() && (strategy.PreferIP4() || strategy.FallbackIP4())), IPv6Enable: (localAddr == nil && strategy.PreferIP6()) || (localAddr != nil && localAddr.Family().IsIPv6() && (strategy.PreferIP6() || strategy.FallbackIP6())), }) { // Resolve fallback if (len(ips) == 0 || err != nil) && strategy.HasFallback() && localAddr == nil { ips, _, err = dnsClient.LookupIP(domain, dns.IPOption{ IPv4Enable: strategy.FallbackIP4(), IPv6Enable: strategy.FallbackIP6(), }) } } if err == nil && len(ips) == 0 { return nil, dns.ErrEmptyResponse } return ips, err } func redirect(ctx context.Context, dst net.Destination, obt string, h outbound.Handler) net.Conn { errors.LogInfo(ctx, "redirecting request "+dst.String()+" to "+obt) outbounds := session.OutboundsFromContext(ctx) ctx = session.ContextWithOutbounds(ctx, append(outbounds, &session.Outbound{ Target: dst, Gateway: nil, Tag: obt, })) // add another outbound in session ctx ur, uw := pipe.New(pipe.OptionsFromContext(ctx)...) dr, dw := pipe.New(pipe.OptionsFromContext(ctx)...) go h.Dispatch(context.WithoutCancel(ctx), &transport.Link{Reader: ur, Writer: dw}) var readerOpt cnc.ConnectionOption if dst.Network == net.Network_TCP { readerOpt = cnc.ConnectionOutputMulti(dr) } else { readerOpt = cnc.ConnectionOutputMultiUDP(dr) } nc := cnc.NewConnection( cnc.ConnectionInputMulti(uw), readerOpt, cnc.ConnectionOnClose(common.ChainedClosable{uw, dw}), ) return nc } func checkAddressPortStrategy(ctx context.Context, dest net.Destination, sockopt *SocketConfig) (*net.Destination, error) { if sockopt.AddressPortStrategy == AddressPortStrategy_None { return nil, nil } newDest := dest var OverridePort, OverrideAddress bool var OverrideBy string switch sockopt.AddressPortStrategy { case AddressPortStrategy_SrvPortOnly: OverridePort = true OverrideAddress = false OverrideBy = "srv" case AddressPortStrategy_SrvAddressOnly: OverridePort = false OverrideAddress = true OverrideBy = "srv" case AddressPortStrategy_SrvPortAndAddress: OverridePort = true OverrideAddress = true OverrideBy = "srv" case AddressPortStrategy_TxtPortOnly: OverridePort = true OverrideAddress = false OverrideBy = "txt" case AddressPortStrategy_TxtAddressOnly: OverridePort = false OverrideAddress = true OverrideBy = "txt" case AddressPortStrategy_TxtPortAndAddress: OverridePort = true OverrideAddress = true OverrideBy = "txt" default: return nil, errors.New("unknown AddressPortStrategy") } if !dest.Address.Family().IsDomain() { return nil, nil } if OverrideBy == "srv" { errors.LogDebug(ctx, "query SRV record for "+dest.Address.String()) parts := strings.SplitN(dest.Address.String(), ".", 3) if len(parts) != 3 { return nil, errors.New("invalid address format", dest.Address.String()) } _, srvRecords, err := net.DefaultResolver.LookupSRV(context.Background(), parts[0][1:], parts[1][1:], parts[2]) if err != nil { return nil, errors.New("failed to lookup SRV record").Base(err) } errors.LogDebug(ctx, "SRV record: "+fmt.Sprintf("addr=%s, port=%d, priority=%d, weight=%d", srvRecords[0].Target, srvRecords[0].Port, srvRecords[0].Priority, srvRecords[0].Weight)) if OverridePort { newDest.Port = net.Port(srvRecords[0].Port) } if OverrideAddress { newDest.Address = net.ParseAddress(srvRecords[0].Target) } return &newDest, nil } if OverrideBy == "txt" { errors.LogDebug(ctx, "query TXT record for "+dest.Address.String()) txtRecords, err := net.DefaultResolver.LookupTXT(ctx, dest.Address.String()) if err != nil { errors.LogError(ctx, "failed to lookup SRV record: "+err.Error()) return nil, errors.New("failed to lookup SRV record").Base(err) } for _, txtRecord := range txtRecords { errors.LogDebug(ctx, "TXT record: "+txtRecord) addr_s, port_s, _ := net.SplitHostPort(string(txtRecord)) addr := net.ParseAddress(addr_s) port, err := net.PortFromString(port_s) if err != nil { continue } if OverridePort { newDest.Port = port } if OverrideAddress { newDest.Address = addr } return &newDest, nil } } return nil, nil } // DialSystem calls system dialer to create a network connection. func DialSystem(ctx context.Context, dest net.Destination, sockopt *SocketConfig) (net.Conn, error) { var src net.Address outbounds := session.OutboundsFromContext(ctx) var outboundName string var origTargetAddr net.Address if len(outbounds) > 0 { ob := outbounds[len(outbounds)-1] if sockopt == nil || len(sockopt.DialerProxy) == 0 { src = ob.Gateway } outboundName = ob.Name origTargetAddr = ob.OriginalTarget.Address if origTargetAddr == nil { origTargetAddr = ob.Target.Address } } if sockopt == nil { return effectiveSystemDialer.Dial(ctx, src, dest, sockopt) } if newDest, err := checkAddressPortStrategy(ctx, dest, sockopt); err == nil && newDest != nil { errors.LogInfo(ctx, "replace destination with "+newDest.String()) dest = *newDest } if sockopt.DomainStrategy.HasStrategy() && dest.Address.Family().IsDomain() { finalStrategy := sockopt.DomainStrategy if outboundName == "freedom" && dest.Network == net.Network_UDP && origTargetAddr != nil && src == nil { finalStrategy = finalStrategy.GetDynamicStrategy(origTargetAddr.Family()) } ips, err := LookupForIP(dest.Address.Domain(), finalStrategy, src) if err != nil { errors.LogErrorInner(ctx, err, "failed to resolve ip") if sockopt.DomainStrategy.ForceIP() { return nil, err } } else if sockopt.HappyEyeballs == nil || sockopt.HappyEyeballs.TryDelayMs == 0 || sockopt.HappyEyeballs.MaxConcurrentTry == 0 || len(ips) < 2 || len(sockopt.DialerProxy) > 0 || dest.Network != net.Network_TCP { dest.Address = net.IPAddress(ips[dice.Roll(len(ips))]) errors.LogInfo(ctx, "replace destination with "+dest.String()) } else { return TcpRaceDial(ctx, src, ips, dest.Port, sockopt, dest.Address.String()) } } if len(sockopt.DialerProxy) > 0 { var obm outbound.Manager if state := systemDialer.Load(); state != nil { obm = state.obm } if obm == nil { return nil, errors.New("there is no outbound manager for dialerProxy").AtError() } h := obm.GetHandler(sockopt.DialerProxy) if h == nil { return nil, errors.New("there is no outbound handler for dialerProxy").AtError() } return redirect(ctx, dest, sockopt.DialerProxy, h), nil } return effectiveSystemDialer.Dial(ctx, src, dest, sockopt) } func InitSystemDialer(dc dns.Client, om outbound.Manager) { for { current := systemDialer.Load() if current != nil && current.held { return } if systemDialer.CompareAndSwap(current, &systemDialerState{dnsClient: dc, obm: om}) { return } } } // HoldSystemDialer (Spanghew): the dialer works with dc and om until ReleaseSystemDialer, whatever // instances are created meanwhile. func HoldSystemDialer(dc dns.Client, om outbound.Manager) { systemDialer.Store(&systemDialerState{dnsClient: dc, obm: om, held: true}) } // ReleaseSystemDialer (Spanghew): the next InitSystemDialer takes effect again; until then the // dialer keeps what it was held with. func ReleaseSystemDialer() { for { current := systemDialer.Load() if current == nil || !current.held { return } if systemDialer.CompareAndSwap(current, &systemDialerState{dnsClient: current.dnsClient, obm: current.obm}) { return } } } ============================================================================== transport/internet/spanghew_system_dialer_test.go (изменённый файл целиком) package internet import ( "context" "sync" "testing" "github.com/xtls/xray-core/common/net" "github.com/xtls/xray-core/features/dns" "github.com/xtls/xray-core/features/outbound" "github.com/xtls/xray-core/transport" ) // Two DNS clients of different types, as the running core (app/dns) and a measurement instance // (the default local client) have: a reader must never see a mix of the two. type spanghewDNSOne struct{ ip net.IP } func (*spanghewDNSOne) Type() interface{} { return dns.ClientType() } func (*spanghewDNSOne) Start() error { return nil } func (*spanghewDNSOne) Close() error { return nil } func (c *spanghewDNSOne) LookupIP(string, dns.IPOption) ([]net.IP, uint32, error) { return []net.IP{c.ip}, 1, nil } type spanghewDNSTwo struct { answers map[string]net.IP } func (spanghewDNSTwo) Type() interface{} { return dns.ClientType() } func (spanghewDNSTwo) Start() error { return nil } func (spanghewDNSTwo) Close() error { return nil } func (c spanghewDNSTwo) LookupIP(domain string, _ dns.IPOption) ([]net.IP, uint32, error) { return []net.IP{c.answers[domain]}, 1, nil } type spanghewOutbounds struct { outbound.Manager handler outbound.Handler } func (m *spanghewOutbounds) GetHandler(string) outbound.Handler { return m.handler } type spanghewOutbound struct{ outbound.Handler } func (spanghewOutbound) Dispatch(context.Context, *transport.Link) {} func spanghewSaveSystemDialer(t *testing.T) { t.Helper() saved := systemDialer.Load() t.Cleanup(func() { systemDialer.Store(saved) }) } func spanghewLookup(t *testing.T) string { t.Helper() ips, err := LookupForIP("example.invalid", DomainStrategy_USE_IP, nil) if err != nil || len(ips) != 1 { t.Fatalf("lookup: %v %v", ips, err) } return ips[0].String() } // While the dialer is held (by the running core), the InitSystemDialer of another instance (a // measurement next to it) changes neither the DNS client nor the outbound manager. func TestSpanghewSystemDialerHeld(t *testing.T) { spanghewSaveSystemDialer(t) running := &spanghewDNSOne{ip: net.IP{203, 0, 113, 1}} measuring := spanghewDNSTwo{answers: map[string]net.IP{"example.invalid": {203, 0, 113, 2}}} proxied := &spanghewOutbounds{handler: spanghewOutbound{}} dialThroughProxy := func() error { conn, err := DialSystem(context.Background(), net.TCPDestination(net.DomainAddress("example.invalid"), 443), &SocketConfig{DialerProxy: "next"}) if err == nil { conn.Close() } return err } InitSystemDialer(measuring, nil) if got := spanghewLookup(t); got != "203.0.113.2" { t.Fatalf("not held: lookup through %s", got) } if dialThroughProxy() == nil { t.Fatal("not held: dialerProxy found an outbound without a manager") } HoldSystemDialer(running, proxied) InitSystemDialer(measuring, nil) if got := spanghewLookup(t); got != "203.0.113.1" { t.Fatalf("held: another instance took the DNS client, lookup through %s", got) } if err := dialThroughProxy(); err != nil { t.Fatalf("held: another instance took the outbound manager: %v", err) } // Released: the dialer keeps what it had until the next instance is created. ReleaseSystemDialer() if got := spanghewLookup(t); got != "203.0.113.1" { t.Fatalf("released: lookup through %s", got) } InitSystemDialer(measuring, nil) if got := spanghewLookup(t); got != "203.0.113.2" { t.Fatalf("after release: lookup through %s", got) } ReleaseSystemDialer() // not held: nothing changes if got := spanghewLookup(t); got != "203.0.113.2" { t.Fatalf("release of a dialer that is not held: lookup through %s", got) } } // Instances come and go while connections resolve names: every lookup gets the answer of one whole // client (run with -race to see the readers and the writer). func TestSpanghewSystemDialerSwappedUnderLookups(t *testing.T) { spanghewSaveSystemDialer(t) one := &spanghewDNSOne{ip: net.IP{203, 0, 113, 1}} two := spanghewDNSTwo{answers: map[string]net.IP{"example.invalid": {203, 0, 113, 2}}} InitSystemDialer(one, nil) stop := make(chan struct{}) var readers sync.WaitGroup for range 4 { readers.Add(1) go func() { defer readers.Done() for { select { case <-stop: return default: } ips, err := LookupForIP("example.invalid", DomainStrategy_USE_IP, nil) if err != nil || len(ips) != 1 || (ips[0].String() != "203.0.113.1" && ips[0].String() != "203.0.113.2") { t.Errorf("lookup: %v %v", ips, err) return } } }() } for i := range 20000 { if i%2 == 0 { InitSystemDialer(two, nil) } else { InitSystemDialer(one, nil) } } close(stop) readers.Wait() } ============================================================================== proxy/tun/stack_gvisor.go (изменённый файл целиком) package tun import ( "context" "time" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/net" "gvisor.dev/gvisor/pkg/buffer" "gvisor.dev/gvisor/pkg/tcpip" "gvisor.dev/gvisor/pkg/tcpip/adapters/gonet" "gvisor.dev/gvisor/pkg/tcpip/checksum" "gvisor.dev/gvisor/pkg/tcpip/header" "gvisor.dev/gvisor/pkg/tcpip/network/ipv4" "gvisor.dev/gvisor/pkg/tcpip/network/ipv6" "gvisor.dev/gvisor/pkg/tcpip/stack" "gvisor.dev/gvisor/pkg/tcpip/transport/icmp" "gvisor.dev/gvisor/pkg/tcpip/transport/tcp" "gvisor.dev/gvisor/pkg/tcpip/transport/udp" "gvisor.dev/gvisor/pkg/waiter" ) const ( defaultNIC tcpip.NICID = 1 tcpRXBufMinSize = tcp.MinBufferSize tcpRXBufDefSize = tcp.DefaultSendBufferSize tcpTXBufMinSize = tcp.MinBufferSize tcpTXBufDefSize = tcp.DefaultReceiveBufferSize ) // stackGVisor is ip stack implemented by gVisor package type stackGVisor struct { ctx context.Context tun Tun idleTimeout time.Duration handler *Handler stack *stack.Stack endpoint stack.LinkEndpoint } // NewStack builds new ip stack (using gVisor) func NewStack(ctx context.Context, options StackOptions, handler *Handler) (Stack, error) { gStack := &stackGVisor{ ctx: ctx, tun: options.Tun, idleTimeout: options.IdleTimeout, handler: handler, } return gStack, nil } // Start is called by Handler to bring stack to life func (t *stackGVisor) Start() error { linkEndpoint, err := t.tun.newEndpoint() if err != nil { return err } ipStack, err := createStack(linkEndpoint) if err != nil { return err } tcpForwarder := tcp.NewForwarder(ipStack, 0, 65535, func(r *tcp.ForwarderRequest) { go func(r *tcp.ForwarderRequest) { var wq waiter.Queue id := r.ID() // Perform a TCP three-way handshake. ep, err := r.CreateEndpoint(&wq) if err != nil { // Spanghew: the app did not finish the handshake of its own connection (it reset or // dropped it) - an everyday event, not a fault of the core. Upstream logged it as an // error: dozens of "connection reset by peer" lines a day in the log the user sees, // pushing the lines that matter out of it. errors.LogInfo(t.ctx, "connection from the tunnel not established: ", err.String()) r.Complete(true) return } options := ep.SocketOptions() options.SetKeepAlive(false) options.SetReuseAddress(true) options.SetReusePort(true) t.handler.HandleConnection( gonet.NewTCPConn(&wq, ep), // local address on the gVisor side is connection destination net.TCPDestination(net.IPAddress(id.LocalAddress.AsSlice()), net.Port(id.LocalPort)), ) // close the socket ep.Close() // send connection complete upstream r.Complete(false) }(r) }) ipStack.SetTransportProtocolHandler(tcp.ProtocolNumber, tcpForwarder.HandlePacket) // Use custom UDP packet handler, instead of strict gVisor forwarder, for FullCone NAT support udpForwarder := newUdpConnectionHandler(t.handler.HandleConnection, t.writeRawUDPPacket) ipStack.SetTransportProtocolHandler(udp.ProtocolNumber, func(id stack.TransportEndpointID, pkt *stack.PacketBuffer) bool { data := pkt.Clone().Data().AsRange().ToSlice() // if len(data) == 0 { // return false // } // source/destination of the packet we process as incoming, on gVisor side are Remote/Local // in other terms, src is the side behind tun, dst is the side behind gVisor // this function handle packets passing from the tun to the gVisor, therefore the src/dst assignement srcIP := net.IPAddress(id.RemoteAddress.AsSlice()) dstIP := net.IPAddress(id.LocalAddress.AsSlice()) if srcIP == nil || dstIP == nil { panic(id) } src := net.UDPDestination(srcIP, net.Port(id.RemotePort)) dst := net.UDPDestination(dstIP, net.Port(id.LocalPort)) udpForwarder.HandlePacket(src, dst, data) return true }) ipStack.SetTransportProtocolHandler(icmp.ProtocolNumber4, t.handleICMPv4Packet) ipStack.SetTransportProtocolHandler(icmp.ProtocolNumber6, t.handleICMPv6Packet) t.stack = ipStack t.endpoint = linkEndpoint return nil } func (t *stackGVisor) writeRawUDPPacket(payload []byte, src net.Destination, dst net.Destination) error { udpLen := header.UDPMinimumSize + len(payload) srcIP := tcpip.AddrFromSlice(src.Address.IP()) dstIP := tcpip.AddrFromSlice(dst.Address.IP()) // build packet with appropriate IP header size isIPv4 := dst.Address.Family().IsIPv4() ipHdrSize := header.IPv6MinimumSize ipProtocol := header.IPv6ProtocolNumber if isIPv4 { ipHdrSize = header.IPv4MinimumSize ipProtocol = header.IPv4ProtocolNumber } pkt := stack.NewPacketBuffer(stack.PacketBufferOptions{ ReserveHeaderBytes: ipHdrSize + header.UDPMinimumSize, Payload: buffer.MakeWithData(payload), }) defer pkt.DecRef() // Build UDP header udpHdr := header.UDP(pkt.TransportHeader().Push(header.UDPMinimumSize)) udpHdr.Encode(&header.UDPFields{ SrcPort: uint16(src.Port), DstPort: uint16(dst.Port), Length: uint16(udpLen), }) // Calculate and set UDP checksum xsum := header.PseudoHeaderChecksum(header.UDPProtocolNumber, srcIP, dstIP, uint16(udpLen)) udpHdr.SetChecksum(^udpHdr.CalculateChecksum(checksum.Checksum(payload, xsum))) // Build IP header if isIPv4 { ipHdr := header.IPv4(pkt.NetworkHeader().Push(header.IPv4MinimumSize)) ipHdr.Encode(&header.IPv4Fields{ TotalLength: uint16(header.IPv4MinimumSize + udpLen), TTL: 64, Protocol: uint8(header.UDPProtocolNumber), SrcAddr: srcIP, DstAddr: dstIP, }) ipHdr.SetChecksum(^ipHdr.CalculateChecksum()) } else { ipHdr := header.IPv6(pkt.NetworkHeader().Push(header.IPv6MinimumSize)) ipHdr.Encode(&header.IPv6Fields{ PayloadLength: uint16(udpLen), TransportProtocol: header.UDPProtocolNumber, HopLimit: 64, SrcAddr: srcIP, DstAddr: dstIP, }) } // dispatch the packet err := t.stack.WriteRawPacket(defaultNIC, ipProtocol, buffer.MakeWithView(pkt.ToView())) if err != nil { return errors.New("failed to write raw udp packet back to stack", err) } return nil } // Close is called by Handler to shut down the stack func (t *stackGVisor) Close() error { if t.stack == nil { return nil } t.endpoint.Attach(nil) t.stack.Close() for _, endpoint := range t.stack.CleanupEndpoints() { endpoint.Abort() } return nil } // createStack configure gVisor ip stack func createStack(ep stack.LinkEndpoint) (*stack.Stack, error) { opts := stack.Options{ NetworkProtocols: []stack.NetworkProtocolFactory{ipv4.NewProtocol, ipv6.NewProtocol}, TransportProtocols: []stack.TransportProtocolFactory{tcp.NewProtocol, udp.NewProtocol, icmp.NewProtocol4, icmp.NewProtocol6}, HandleLocal: false, } gStack := stack.New(opts) err := gStack.CreateNIC(defaultNIC, ep) if err != nil { return nil, errors.New(err.String()) } gStack.SetRouteTable([]tcpip.Route{ {Destination: header.IPv4EmptySubnet, NIC: defaultNIC}, {Destination: header.IPv6EmptySubnet, NIC: defaultNIC}, }) err = gStack.SetSpoofing(defaultNIC, true) if err != nil { return nil, errors.New(err.String()) } err = gStack.SetPromiscuousMode(defaultNIC, true) if err != nil { return nil, errors.New(err.String()) } cOpt := tcpip.CongestionControlOption("cubic") gStack.SetTransportProtocolOption(tcp.ProtocolNumber, &cOpt) sOpt := tcpip.TCPSACKEnabled(true) gStack.SetTransportProtocolOption(tcp.ProtocolNumber, &sOpt) mOpt := tcpip.TCPModerateReceiveBufferOption(true) gStack.SetTransportProtocolOption(tcp.ProtocolNumber, &mOpt) // Disable RACK/TLP loss recovery to fix connection stalls under high load rOpt := tcpip.TCPRecovery(0) gStack.SetTransportProtocolOption(tcp.ProtocolNumber, &rOpt) tcpRXBufOpt := tcpip.TCPReceiveBufferSizeRangeOption{ Min: tcpRXBufMinSize, Default: tcpRXBufDefSize, Max: tcpRXBufMaxSize, } err = gStack.SetTransportProtocolOption(tcp.ProtocolNumber, &tcpRXBufOpt) if err != nil { return nil, errors.New(err.String()) } tcpTXBufOpt := tcpip.TCPSendBufferSizeRangeOption{ Min: tcpTXBufMinSize, Default: tcpTXBufDefSize, Max: tcpTXBufMaxSize, } err = gStack.SetTransportProtocolOption(tcp.ProtocolNumber, &tcpTXBufOpt) if err != nil { return nil, errors.New(err.String()) } return gStack, nil } ============================================================================== proxy/tun/handler.go (изменённый файл целиком) package tun import ( "context" "net/netip" "runtime" "strings" "sync" "syscall" "github.com/xtls/xray-core/common" "github.com/xtls/xray-core/common/buf" c "github.com/xtls/xray-core/common/ctx" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/log" "github.com/xtls/xray-core/common/net" "github.com/xtls/xray-core/common/protocol" "github.com/xtls/xray-core/common/session" "github.com/xtls/xray-core/core" "github.com/xtls/xray-core/features/policy" "github.com/xtls/xray-core/features/routing" "github.com/xtls/xray-core/features/stats" "github.com/xtls/xray-core/transport" "github.com/xtls/xray-core/transport/internet" "github.com/xtls/xray-core/transport/internet/stat" ) // Handler is managing object that tie together tun interface, ip stack and dispatch connections to the routing type Handler struct { ctx context.Context config *Config stack Stack tun Tun policyManager policy.Manager dispatcher routing.Dispatcher tag string sniffingRequest session.SniffingRequest uplinkCounter stats.Counter downlinkCounter stats.Counter } // ConnectionHandler interface with the only method that stack is going to push new connections to type ConnectionHandler interface { HandleConnection(conn net.Conn, destination net.Destination) } // Handler implements ConnectionHandler var _ ConnectionHandler = (*Handler)(nil) // Handler implements common.Runnable var _ common.Runnable = (*Handler)(nil) // Init the Handler instance with necessary parameters func (t *Handler) Init(ctx context.Context, pm policy.Manager, dispatcher routing.Dispatcher) error { // Retrieve tag and sniffing config from context (set by AlwaysOnInboundHandler) if inbound := session.InboundFromContext(ctx); inbound != nil { t.tag = inbound.Tag } if content := session.ContentFromContext(ctx); content != nil { t.sniffingRequest = content.SniffingRequest } t.ctx = core.ToBackgroundDetachedContext(ctx) t.policyManager = pm t.dispatcher = dispatcher if len(t.tag) > 0 && pm.ForSystem().Stats.InboundUplink { statsManager := core.MustFromContext(ctx).GetFeature(stats.ManagerType()).(stats.Manager) name := "inbound>>>" + t.tag + ">>>traffic>>>uplink" c, _ := statsManager.GetOrRegisterCounter(name) if c != nil { t.uplinkCounter = c } } if len(t.tag) > 0 && pm.ForSystem().Stats.InboundDownlink { statsManager := core.MustFromContext(ctx).GetFeature(stats.ManagerType()).(stats.Manager) name := "inbound>>>" + t.tag + ">>>traffic>>>downlink" c, _ := statsManager.GetOrRegisterCounter(name) if c != nil { t.downlinkCounter = c } } return nil } // Spanghew: see Start. var interfaceControllerOnce sync.Once func bindToOutboundInterface(network, address string, c syscall.RawConn) error { addrPort, _ := netip.ParseAddrPort(address) // skip loopback if addrPort.Addr().IsLoopback() || strings.HasPrefix(strings.ToLower(address), "localhost:") { return nil } iface := updater.Get() if iface == nil { // Spanghew: no interface to bind to (no network besides the TUN). Upstream left // the socket unbound: it followed the default route back into the TUN, and the // core talked to itself in a loop. No network - the dial fails: the dialer stops // on this error only (transport/internet/system_dialer.go), any other it logs. return internet.ErrControllerRefused } return c.Control(func(fd uintptr) { err := setinterface(network, address, fd, iface) if err != nil { errors.LogInfoInner(context.Background(), err, "[tun] falied to set interface") } }) } func (t *Handler) Start() error { tunName := t.config.Name tunInterface, err := NewTun(t.config) if err != nil { return err } if t.config.AutoOutboundsInterface != "" { tunIndex, err := tunInterface.Index() if err != nil { _ = tunInterface.Close() return err } if t.config.AutoOutboundsInterface == "auto" { t.config.AutoOutboundsInterface = "" } updater = &InterfaceUpdater{tunIndex: tunIndex, fixedName: t.config.AutoOutboundsInterface} updater.Update() // Spanghew: one controller per process. It reads the package-level updater, and // upstream appended another copy on every start of the inbound - a long-lived // process (the Windows service) bound each socket once per past connection. interfaceControllerOnce.Do(func() { internet.RegisterDialerController(bindToOutboundInterface) }) } errors.LogInfo(t.ctx, tunName, " created") tunStackOptions := StackOptions{ Tun: tunInterface, IdleTimeout: t.policyManager.ForLevel(t.config.UserLevel).Timeouts.ConnectionIdle, } tunStack, err := NewStack(t.ctx, tunStackOptions, t) if err != nil { _ = tunInterface.Close() return err } err = tunStack.Start() if err != nil { _ = tunStack.Close() _ = tunInterface.Close() return err } err = tunInterface.Start() if err != nil { _ = tunStack.Close() _ = tunInterface.Close() return err } t.stack = tunStack t.tun = tunInterface errors.LogInfo(t.ctx, tunName, " up") return nil } // HandleConnection pass the connection coming from the ip stack to the routing dispatcher func (t *Handler) HandleConnection(conn net.Conn, destination net.Destination) { // when handling is done with any outcome, always signal back to the incoming connection // to close, send completion packets back to the network, and cleanup defer conn.Close() ctx, cancel := context.WithCancel(t.ctx) defer cancel() ctx = c.ContextWithID(ctx, session.NewID()) // if the connection is already closed, conn.RemoteAddr() will be nil // due to gvisor weird behavior remote := conn.RemoteAddr() if remote == nil { errors.LogInfo(t.ctx, "dropped quickly closed connection") return } source := net.DestinationFromAddr(remote) if !ownerAllowed(source, destination) { errors.LogInfo(t.ctx, "blocked connection of an excluded or unknown app from ", source, " to ", destination) return } if t.uplinkCounter != nil || t.downlinkCounter != nil { conn = &stat.CounterConnection{ Connection: conn, ReadCounter: t.uplinkCounter, WriteCounter: t.downlinkCounter, } } inbound := session.Inbound{ Name: "tun", Tag: t.tag, CanSpliceCopy: 3, Source: source, User: &protocol.MemoryUser{ Level: t.config.UserLevel, }, } ctx = session.ContextWithInbound(ctx, &inbound) ctx = session.ContextWithContent(ctx, &session.Content{ SniffingRequest: t.sniffingRequest, }) ctx = session.SubContextFromMuxInbound(ctx) ctx = log.ContextWithAccessMessage(ctx, &log.AccessMessage{ From: inbound.Source, To: destination, Status: log.AccessAccepted, Reason: "", }) errors.LogInfo(ctx, "processing from ", source, " to ", destination) link := &transport.Link{ Reader: &buf.TimeoutWrapperReader{Reader: buf.NewReader(conn)}, Writer: buf.NewWriter(conn), } if err := t.dispatcher.DispatchLink(ctx, destination, link); err != nil { errors.LogError(ctx, errors.New("connection closed").Base(err)) } } // Close implements common.Closable. func (t *Handler) Close() error { return errors.Combine(common.CloseIfExists(t.stack), common.CloseIfExists(t.tun)) } // Network implements proxy.Inbound // and exists only to comply to proxy interface, declaring it doesn't listen on any network, // making the process not open any port for this inbound (input will be network interface) func (t *Handler) Network() []net.Network { return []net.Network{} } // Process implements proxy.Inbound // and exists only to comply to proxy interface, which should never get any inputs due to no listening ports func (t *Handler) Process(ctx context.Context, network net.Network, conn stat.Connection, dispatcher routing.Dispatcher) error { return nil } func init() { common.Must(common.RegisterConfig((*Config)(nil), func(ctx context.Context, config interface{}) (interface{}, error) { t := &Handler{config: config.(*Config)} err := core.RequireFeatures(ctx, func(pm policy.Manager, dispatcher routing.Dispatcher) error { return t.Init(ctx, pm, dispatcher) }) return t, err })) } // ownerAllowed (Spanghew): пускать ли соединение в туннель, по владельцу — uid приложения // от net.FindProcess (на Android его ищет приложение через libXray.RegisterProcessFinder и // ConnectivityManager.getConnectionOwnerUid). Отрицательный uid — «не пускать»: владелец не найден // (так бывает у сокета, привязанного к tun0 через SO_BINDTODEVICE, — приложение из исключений VPN // так узнало бы адрес сервера) или приложение исключено из VPN. Поиск не зарегистрирован или не // смог спросить — пропускаем, как без заплатки. Адреса — исходные, до sniffing. // Только Android: на macOS и Windows FindProcess отдаёт PID ≥ 0 или ошибку (соединение всё равно пропускается), // но на каждое соединение TUN перебирает сокеты всех процессов или таблицы соединений Windows (аудит M24d, w7). func ownerAllowed(source, destination net.Destination) bool { if runtime.GOOS != "android" { return true } if !source.Address.Family().IsIP() || !destination.Address.Family().IsIP() { return true } network := "tcp" if destination.Network == net.Network_UDP { network = "udp" } uid, _, _, err := net.FindProcess(network, source.Address.IP().String(), uint16(source.Port), destination.Address.IP().String(), uint16(destination.Port)) return err != nil || uid >= 0 } ============================================================================== proxy/tun/udp_fullcone.go (изменённый файл целиком) package tun import ( "context" "io" "net/netip" "sync" "time" "github.com/xtls/xray-core/common/buf" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/net" ) type packet struct { data []byte dest *net.Destination } // sub-handler specifically for udp connections under main handler type udpConnectionHandler struct { sync.RWMutex udpConns map[udpConnKey]*udpConn handleConnection func(conn net.Conn, dest net.Destination) writePacket func(data []byte, src net.Destination, dst net.Destination) error } func newUdpConnectionHandler(handleConnection func(conn net.Conn, dest net.Destination), writePacket func(data []byte, src net.Destination, dst net.Destination) error) *udpConnectionHandler { handler := &udpConnectionHandler{ udpConns: make(map[udpConnKey]*udpConn), handleConnection: handleConnection, writePacket: writePacket, } return handler } // Spanghew: the route of a connection is chosen by its first packet. Upstream keyed connections by // the source alone, so a socket whose first packet went to a local address (routed "direct" or to // the local-network block) carried all its later packets to public addresses the same way - past // the VPN. Local and public destinations of one socket are separate connections now; FullCone NAT // stays within each. type udpConnKey struct { src net.Destination local bool } func udpKey(src, dst net.Destination) udpConnKey { return udpConnKey{src: src, local: isLocalDestination(dst)} } // localNetworks: geoip:private of the app's geobase (v2fly/geoip) - what subscription templates // route "direct". Private (RFC 1918, fc00::/7), loopback, link-local, multicast and broadcast // addresses are in it, and so are the reserved ranges a public host never answers from // (0.0.0.0/8, 100.64.0.0/10, 192.0.2.0/24, 198.18.0.0/15, 240.0.0.0/4 ...): a first packet there // would take the "direct" route as well. The same list - udp_pretend_scope of the hev lwIP patch; // the app's test compares both with the geobase. var localNetworks = func() []netip.Prefix { var networks []netip.Prefix for _, network := range []string{ "0.0.0.0/8", "10.0.0.0/8", "100.64.0.0/10", "127.0.0.0/8", "169.254.0.0/16", "172.16.0.0/12", "192.0.0.0/24", "192.0.2.0/24", "192.88.99.0/24", "192.168.0.0/16", "198.18.0.0/15", "198.51.100.0/24", "203.0.113.0/24", "224.0.0.0/3", "::/127", "fc00::/7", "fe80::/10", "ff00::/8", } { networks = append(networks, netip.MustParsePrefix(network)) } return networks }() // isLocalDestination: an address of localNetworks (IPv4 inside IPv6 "::ffff:..." - as IPv4, the way // the router reads it). func isLocalDestination(dst net.Destination) bool { if dst.Address == nil || !dst.Address.Family().IsIP() { return false } ip, ok := netip.AddrFromSlice(dst.Address.IP()) if !ok { return false } ip = ip.Unmap() for _, network := range localNetworks { if network.Contains(ip) { return true } } return false } // HandlePacket handles UDP packets coming from tun, to forward to the dispatcher // this custom handler support FullCone NAT of returning packets, binding connection only by the source addr:port func (u *udpConnectionHandler) HandlePacket(src net.Destination, dst net.Destination, data []byte) { key := udpKey(src, dst) u.RLock() conn, found := u.udpConns[key] if found { select { case conn.egress <- &packet{ data: data, dest: &dst, }: default: errors.LogDebug(context.Background(), "drop udp with size ", len(data), " to ", dst.NetAddr(), " original ", conn.dst.NetAddr(), " > queue full") } u.RUnlock() return } u.RUnlock() u.Lock() defer u.Unlock() conn, found = u.udpConns[key] if !found { egress := make(chan *packet, 1024) conn = &udpConn{handler: u, egress: egress, src: src, dst: dst} u.udpConns[key] = conn go u.handleConnection(conn, dst) } // send packet data to the egress channel, if it has buffer, or discard select { case conn.egress <- &packet{ data: data, dest: &dst, }: default: errors.LogDebug(context.Background(), "drop udp with size ", len(data), " to ", dst.NetAddr(), " original ", conn.dst.NetAddr(), " > queue full 2") } } func (u *udpConnectionHandler) connectionFinished(key udpConnKey) { u.Lock() conn, found := u.udpConns[key] if found { delete(u.udpConns, key) close(conn.egress) } u.Unlock() } // udp connection abstraction type udpConn struct { handler *udpConnectionHandler egress chan *packet src net.Destination dst net.Destination } func (c *udpConn) ReadMultiBuffer() (buf.MultiBuffer, error) { for { e, ok := <-c.egress if !ok { return nil, io.EOF } b := buf.New() if _, err := b.Write(e.data); err != nil { errors.LogErrorInner(context.Background(), err, "drop packet to ", e.dest, " with size ", len(e.data)) b.Release() continue } b.UDP = e.dest return buf.MultiBuffer{b}, nil } } // Read packets from the connection func (c *udpConn) Read(p []byte) (int, error) { e, ok := <-c.egress if !ok { return 0, io.EOF } n := copy(p, e.data) if n != len(e.data) { return 0, io.ErrShortBuffer } return n, nil } func (c *udpConn) WriteMultiBuffer(mb buf.MultiBuffer) error { for i, b := range mb { dst := c.dst if b.UDP != nil { if b.UDP.Address.Family().IsDomain() { errors.LogError(context.Background(), "impossible domain packet ", b.UDP, " reply via original target ", dst) } else { dst = *b.UDP } } err := c.handler.writePacket(b.Bytes(), dst, c.src) if err != nil { buf.ReleaseMulti(mb[i:]) return err } b.Release() } return nil } // Write returning packets back func (c *udpConn) Write(p []byte) (int, error) { // sending packets back mean sending payload with source/destination reversed err := c.handler.writePacket(p, c.dst, c.src) if err != nil { return 0, nil } return len(p), nil } func (c *udpConn) Close() error { c.handler.connectionFinished(udpKey(c.src, c.dst)) return nil } func (c *udpConn) LocalAddr() net.Addr { return c.dst.RawNetAddr() } func (c *udpConn) RemoteAddr() net.Addr { return c.src.RawNetAddr() } func (c *udpConn) SetDeadline(t time.Time) error { return nil } func (c *udpConn) SetReadDeadline(t time.Time) error { return nil } func (c *udpConn) SetWriteDeadline(t time.Time) error { return nil } ============================================================================== proxy/tun/spanghew_udp_scope_test.go (изменённый файл целиком) package tun import ( "sync" "testing" "github.com/xtls/xray-core/common/net" ) // Spanghew: one socket, first packet to the local network, the next to a public address - two // connections, each routed by its own first packet; packets to other public addresses stay in the // public one (FullCone). func TestSpanghewUDPLocalAndPublicAreSeparateConnections(t *testing.T) { var ( mu sync.Mutex opened []net.Destination conns []net.Conn wait sync.WaitGroup ) handler := newUdpConnectionHandler(func(conn net.Conn, dest net.Destination) { mu.Lock() opened = append(opened, dest) conns = append(conns, conn) mu.Unlock() wait.Done() }, func([]byte, net.Destination, net.Destination) error { return nil }) src := net.UDPDestination(net.ParseAddress("172.19.0.1"), 50000) router := net.UDPDestination(net.ParseAddress("192.168.1.1"), 3478) public := net.UDPDestination(net.ParseAddress("1.1.1.1"), 3478) other := net.UDPDestination(net.ParseAddress("8.8.8.8"), 3478) wait.Add(2) handler.HandlePacket(src, router, []byte{1}) handler.HandlePacket(src, public, []byte{2}) handler.HandlePacket(src, other, []byte{3}) handler.HandlePacket(src, router, []byte{4}) wait.Wait() mu.Lock() defer mu.Unlock() if len(opened) != 2 { t.Fatalf("connections: %v", opened) } seen := map[net.Destination]bool{opened[0]: true, opened[1]: true} if !seen[router] || !seen[public] { t.Fatalf("connections routed by: %v", opened) } if len(handler.udpConns) != 2 { t.Fatalf("connections kept: %d", len(handler.udpConns)) } for _, conn := range conns { _ = conn.Close() } if len(handler.udpConns) != 0 { t.Fatalf("connections left after close: %d", len(handler.udpConns)) } } func TestSpanghewLocalDestinations(t *testing.T) { for address, local := range map[string]bool{ "10.1.2.3": true, "172.16.0.9": true, "192.168.1.255": true, "169.254.10.10": true, "127.0.0.1": true, "224.0.0.251": true, "255.255.255.255": true, "0.1.2.3": true, "100.64.0.1": true, "100.127.255.255": true, "192.0.0.9": true, "192.0.2.1": true, "192.88.99.1": true, "198.18.0.1": true, "198.19.255.255": true, "198.51.100.7": true, "203.0.113.7": true, "240.0.0.1": true, "::": true, "::1": true, "fe80::1": true, "fd00::1": true, "ff02::fb": true, "::ffff:192.0.2.1": true, "::ffff:192.168.1.1": true, "1.1.1.1": false, "100.63.255.255": false, "100.128.0.1": false, "192.0.1.1": false, "198.17.255.255": false, "198.20.0.1": false, "223.255.255.255": false, "::2": false, "::ffff:1.1.1.1": false, "2606:4700::1111": false, "2001:db8::1": false, } { if got := isLocalDestination(net.UDPDestination(net.ParseAddress(address), 53)); got != local { t.Errorf("%s: local %v, want %v", address, got, local) } } if isLocalDestination(net.UDPDestination(net.ParseAddress("example.com"), 53)) { t.Error("a domain is not a local destination") } } // The reported leak: the first packet of a socket goes to a reserved address that geoip:private holds // and RFC 1918 does not; the socket's packets to public addresses are another connection. func TestSpanghewUDPReservedAndPublicAreSeparateConnections(t *testing.T) { var ( mu sync.Mutex opened []net.Destination conns []net.Conn wait sync.WaitGroup ) handler := newUdpConnectionHandler(func(conn net.Conn, dest net.Destination) { mu.Lock() opened = append(opened, dest) conns = append(conns, conn) mu.Unlock() wait.Done() }, func([]byte, net.Destination, net.Destination) error { return nil }) src := net.UDPDestination(net.ParseAddress("172.19.0.1"), 50000) reserved := net.UDPDestination(net.ParseAddress("192.0.2.1"), 9) public := net.UDPDestination(net.ParseAddress("1.1.1.1"), 3478) wait.Add(2) handler.HandlePacket(src, reserved, []byte{1}) handler.HandlePacket(src, public, []byte{2}) wait.Wait() mu.Lock() defer mu.Unlock() seen := map[net.Destination]bool{opened[0]: true, opened[1]: true} if len(opened) != 2 || !seen[reserved] || !seen[public] { t.Fatalf("connections routed by: %v", opened) } for _, conn := range conns { _ = conn.Close() } } ============================================================================== proxy/tun/tun_windows.go (изменённый файл целиком) //go:build windows package tun import ( "context" "crypto/md5" "encoding/binary" go_errors "errors" "net" "net/netip" "sync" "time" "unsafe" "github.com/xtls/xray-core/common/errors" "golang.org/x/sys/windows" "golang.zx2c4.com/wintun" "golang.zx2c4.com/wireguard/windows/tunnel/winipcfg" "gvisor.dev/gvisor/pkg/buffer" "gvisor.dev/gvisor/pkg/tcpip" "gvisor.dev/gvisor/pkg/tcpip/stack" ) //go:linkname procyield runtime.procyield func procyield(cycles uint32) // WindowsTun is an object that handles tun network interface on Windows // current version is heavily stripped to do nothing more, // then create a network interface, to be provided as endpoint to gVisor ip stack type WindowsTun struct { sync.RWMutex options *Config adapter *wintun.Adapter session wintun.Session readWait windows.Handle luid winipcfg.LUID cbr winipcfg.ChangeCallback cbi winipcfg.ChangeCallback closed bool } // WindowsTun implements Tun var _ Tun = (*WindowsTun)(nil) // WindowsTun implements GVisorDevice var _ GVisorDevice = (*WindowsTun)(nil) // NewTun creates a Wintun interface with the given name. Should a Wintun // interface with the same name exist, it tried to be reused. func NewTun(options *Config) (Tun, error) { // instantiate wintun adapter adapter, err := open(options.Name, options.Desc) if err != nil { return nil, err } // start the interface with ring buffer capacity of 8 MiB session, err := adapter.StartSession(0x800000) if err != nil { _ = adapter.Close() return nil, err } tun := &WindowsTun{ options: options, adapter: adapter, session: session, readWait: session.ReadWaitEvent(), luid: winipcfg.LUID(adapter.LUID()), } return tun, nil } func open(name, desc string) (*wintun.Adapter, error) { // generate a deterministic GUID from the adapter name id := md5.Sum([]byte(name)) guid := (*windows.GUID)(unsafe.Pointer(&id[0])) // try to create adapter anew adapter, err := wintun.CreateAdapter(name, desc, guid) if err == nil { return adapter, nil } return nil, err } func (t *WindowsTun) Start() (err error) { var address4, address6 bool addresses := make([]netip.Prefix, 0, len(t.options.Gateway)) for _, cidr := range t.options.Gateway { prefix := netip.MustParsePrefix(cidr) if prefix.Addr().Is4() { address4 = true } else { address6 = true } addresses = append(addresses, prefix) } dns := make([]netip.Addr, 0, len(t.options.DNS)) for _, ip := range t.options.DNS { dns = append(dns, netip.MustParseAddr(ip)) } var route4, route6 bool routesMap := make(map[winipcfg.RouteData]struct{}) for _, cidr := range t.options.AutoSystemRoutingTable { prefix := netip.MustParsePrefix(cidr) route := winipcfg.RouteData{ Destination: prefix.Masked(), Metric: 0, } if prefix.Addr().Is4() { route4 = true route.NextHop = netip.IPv4Unspecified() } else { route6 = true route.NextHop = netip.IPv6Unspecified() } routesMap[route] = struct{}{} } routesData := make([]*winipcfg.RouteData, 0, len(routesMap)) for route := range routesMap { r := route routesData = append(routesData, &r) } var retryTimes int var firstErr error startOver: if retryTimes > 0 { if retryTimes > 15 { return windows.ERROR_NOT_FOUND } errors.LogErrorInner(context.Background(), firstErr, "Interface configuration failed, retrying attempt ", retryTimes, "/15") time.Sleep(time.Second) } retryTimes++ for _, family := range []winipcfg.AddressFamily{windows.AF_INET, windows.AF_INET6} { if family == windows.AF_INET && route4 || family == windows.AF_INET6 && route6 { err = t.luid.SetRoutesForFamily(family, routesData) if err != nil { firstErr = errors.New("unable to set routes").Base(err) if err == windows.ERROR_NOT_FOUND { goto startOver } return firstErr } } if family == windows.AF_INET && address4 || family == windows.AF_INET6 && address6 { err = t.luid.SetIPAddressesForFamily(family, addresses) if err != nil { firstErr = errors.New("unable to set ips").Base(err) if err == windows.ERROR_NOT_FOUND { goto startOver } return firstErr } } ipif, err := t.luid.IPInterface(family) if err != nil { return err } ipif.RouterDiscoveryBehavior = winipcfg.RouterDiscoveryDisabled ipif.DadTransmits = 0 ipif.ManagedAddressConfigurationSupported = false ipif.OtherStatefulConfigurationSupported = false if family == windows.AF_INET && (address4 || route4) || family == windows.AF_INET6 && (address6 || route6) { ipif.NLMTU = t.options.MTU } if family == windows.AF_INET && route4 || family == windows.AF_INET6 && route6 { ipif.UseAutomaticMetric = false ipif.Metric = 0 } err = ipif.Set() if err != nil { firstErr = errors.New("unable to set metric and MTU").Base(err) if err == windows.ERROR_NOT_FOUND { goto startOver } return firstErr } err = t.luid.SetDNS(family, dns, nil) if err != nil { firstErr = errors.New("unable to set DNS").Base(err) if err == windows.ERROR_NOT_FOUND { goto startOver } return firstErr } } if updater != nil { t.cbr, err = winipcfg.RegisterRouteChangeCallback(func(notificationType winipcfg.MibNotificationType, route *winipcfg.MibIPforwardRow2) { updater.Update() }) if err != nil { return err } t.cbi, err = winipcfg.RegisterInterfaceChangeCallback(func(notificationType winipcfg.MibNotificationType, iface *winipcfg.MibIPInterfaceRow) { updater.Update() }) if err != nil { return err } // Spanghew: the first Update() ran before the TUN had its routes (handler.go), and what // changed between it and the two registrations above nobody reported: a default route // that appeared in that gap (the network coming up while the tunnel starts) left the // updater without an interface, and the controller refused every dial until the next // change of the network. Look once more, now that changes are reported. updater.Update() } return nil } func (t *WindowsTun) Close() error { t.Lock() defer t.Unlock() if t.closed { return nil } t.closed = true if t.cbr != nil { t.cbr.Unregister() } if t.cbi != nil { t.cbi.Unregister() } // Spanghew: туннеля больше нет — и привязки к его прежней сетевой карте тоже. Контроллер сокетов один на // процесс (handler.go) и читает updater, а обновлять его после Close некому: служба Windows живёт дальше, // и замеры задержки шли бы через карту прошлого туннеля (сменили Wi-Fi на кабель — «нет ответа») или не // шли вовсе (туннель остановили без сети). Номер 0 для IP_UNICAST_IF — «без привязки», маршрут выбирает // Windows. Следующий Start ставит свой updater. if updater != nil { updater = &InterfaceUpdater{iface: &net.Interface{}} } if t.luid != 0 { t.luid.FlushRoutes(windows.AF_INET) t.luid.FlushIPAddresses(windows.AF_INET) t.luid.FlushDNS(windows.AF_INET) t.luid.FlushRoutes(windows.AF_INET6) t.luid.FlushIPAddresses(windows.AF_INET6) t.luid.FlushDNS(windows.AF_INET6) } if t.session != (wintun.Session{}) { t.session.End() } if t.adapter != nil { t.adapter.Close() } return nil } func (t *WindowsTun) Name() (string, error) { row, err := t.luid.Interface() if err != nil { return "", err } return row.Alias(), nil } func (t *WindowsTun) Index() (int, error) { row, err := t.luid.Interface() if err != nil { return 0, err } return int(row.InterfaceIndex), nil } // WritePacket implements GVisorDevice method to write one packet to the tun device func (t *WindowsTun) WritePacket(packetBuffer *stack.PacketBuffer) tcpip.Error { t.RLock() defer t.RUnlock() if t.closed { return &tcpip.ErrClosedForSend{} } // request buffer from Wintun packet, err := t.session.AllocateSendPacket(packetBuffer.Size()) if err != nil { return &tcpip.ErrAborted{} } // copy the bytes of slices that compose the packet into the allocated buffer var index int for _, packetElement := range packetBuffer.AsSlices() { index += copy(packet[index:], packetElement) } // signal Wintun to send that buffer as the packet t.session.SendPacket(packet) return nil } // ReadPacket implements GVisorDevice method to read one packet from the tun device // It is expected that the method will not block, rather return ErrQueueEmpty when there is nothing on the line, // which will make the stack call Wait which should implement desired push-back func (t *WindowsTun) ReadPacket() (byte, *stack.PacketBuffer, error) { packet, err := t.session.ReceivePacket() if go_errors.Is(err, windows.ERROR_NO_MORE_ITEMS) { return 0, nil, ErrQueueEmpty } if err != nil { return 0, nil, err } version := packet[0] >> 4 packetBuffer := buffer.MakeWithView(buffer.NewViewWithData(packet)) return version, stack.NewPacketBuffer(stack.PacketBufferOptions{ Payload: packetBuffer, IsForwardedPacket: true, OnRelease: func() { t.session.ReleaseReceivePacket(packet) }, }), nil } func (t *WindowsTun) Wait() { procyield(1) _, _ = windows.WaitForSingleObject(t.readWait, windows.INFINITE) } func (t *WindowsTun) newEndpoint() (stack.LinkEndpoint, error) { return &LinkEndpoint{deviceMTU: t.options.MTU, device: t}, nil } const ( IP_UNICAST_IF = 31 IPV6_UNICAST_IF = 31 ) func setinterface(network, address string, fd uintptr, iface *net.Interface) error { var index [4]byte binary.BigEndian.PutUint32(index[:], uint32(iface.Index)) var err1, err2, err3, err4 error switch network { case "tcp6", "udp6", "ip6": err1 = windows.SetsockoptInt(windows.Handle(fd), windows.IPPROTO_IPV6, IPV6_UNICAST_IF, iface.Index) if network == "udp6" { err2 = windows.SetsockoptInt(windows.Handle(fd), windows.IPPROTO_IPV6, windows.IPV6_MULTICAST_IF, iface.Index) } fallthrough case "tcp4", "udp4", "ip4": err3 = windows.SetsockoptInt(windows.Handle(fd), windows.IPPROTO_IP, IP_UNICAST_IF, *(*int)(unsafe.Pointer(&index[0]))) if network == "udp4" || network == "udp6" { err4 = windows.SetsockoptInt(windows.Handle(fd), windows.IPPROTO_IP, windows.IP_MULTICAST_IF, *(*int)(unsafe.Pointer(&index[0]))) } default: panic(network + " " + address) } return errors.Combine(err1, err2, err3, err4) } func findOutboundInterface(tunIndex int, fixedName string) (*net.Interface, error) { if fixedName != "" { return net.InterfaceByName(fixedName) } // Spanghew: интерфейс, через который система и так ходит в интернет, — с маршрутом по умолчанию // и наименьшей суммой метрик, кроме самого TUN. Ядро с 26.9.8 тоже выбирает по маршрутам, но Wi-Fi // берёт раньше кабеля с меньшей метрикой и не различает IPv4 и IPv6. Нет такого интерфейса — nil: // контроллер сокетов (handler.go) откажет в соединении, а не пустит его обратно в TUN. index := defaultRouteInterfaceIndex(tunIndex) if index == 0 { return nil, nil } return net.InterfaceByIndex(index) } // defaultRouteInterfaceIndex (Spanghew): номер подключённого интерфейса с маршрутом по умолчанию // и наименьшей суммой метрик маршрута и интерфейса (как выбирает сама Windows); 0 — такого нет. // Сначала IPv4 (0.0.0.0/0); нет — IPv6 (::/0): сеть только с IPv6. func defaultRouteInterfaceIndex(tunIndex int) int { for _, family := range []winipcfg.AddressFamily{windows.AF_INET, windows.AF_INET6} { if index := defaultRouteInterfaceIndexOf(family, tunIndex); index != 0 { return index } } return 0 } func defaultRouteInterfaceIndexOf(family winipcfg.AddressFamily, tunIndex int) int { routes, err := winipcfg.GetIPForwardTable2(family) if err != nil { return 0 } best, bestMetric := 0, uint64(0) for _, route := range routes { if route.DestinationPrefix.PrefixLength != 0 || int(route.InterfaceIndex) == tunIndex { continue } ipif, err := route.InterfaceLUID.IPInterface(family) if err != nil || !ipif.Connected { continue } metric := uint64(route.Metric) + uint64(ipif.Metric) if best == 0 || metric < bestMetric { best, bestMetric = int(route.InterfaceIndex), metric } } return best } ============================================================================== transport/internet/system_dialer.go (изменённый файл целиком) package internet import ( "context" goerrors "errors" "sync" "syscall" "time" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/net" "github.com/xtls/xray-core/features/dns" "github.com/xtls/xray-core/features/outbound" ) var ( Controllers []func(network, address string, c syscall.RawConn) error ControllersLock sync.Mutex effectiveSystemDialer SystemDialer = &DefaultSystemDialer{} ) // ErrControllerRefused (Spanghew): a dialer controller returns it to fail the dial. Upstream only // logs what a controller returns and dials on, so a controller could not stop a socket it had not // prepared: on Windows a socket left unbound with no network besides the TUN followed the default // route back into the TUN (the controller is in proxy/tun/handler.go). Any other error of a // controller is logged and ignored, as before. var ErrControllerRefused = goerrors.New("dial refused by a dialer controller") type SystemDialer interface { Dial(ctx context.Context, source net.Address, destination net.Destination, sockopt *SocketConfig) (net.Conn, error) DestIpAddress() net.IP } type DefaultSystemDialer struct { dns dns.Client obm outbound.Manager } func resolveSrcAddr(network net.Network, src net.Address) net.Addr { if src == nil || src == net.AnyIP { return nil } if network == net.Network_TCP { return &net.TCPAddr{ IP: src.IP(), Port: 0, } } return &net.UDPAddr{ IP: src.IP(), Port: 0, } } func (d *DefaultSystemDialer) Dial(ctx context.Context, src net.Address, dest net.Destination, sockopt *SocketConfig) (net.Conn, error) { errors.LogDebug(ctx, "dialing to "+dest.String()) if dest.Network == net.Network_UDP { destAddr, err := net.ResolveUDPAddr("udp", dest.NetAddr()) if err != nil { return nil, err } srcAddr := resolveSrcAddr(net.Network_UDP, src) if srcAddr == nil { // some OS don't support mapped IPv4 dual stack // and need to select 0.0.0.0 or [::] manually based on the destination wildcard := net.AnyIP.IP() if destAddr.IP.To4() == nil { wildcard = net.AnyIPv6.IP() } srcAddr = &net.UDPAddr{ IP: wildcard, Port: 0, } } var lc net.ListenConfig lc.Control = func(network, address string, c syscall.RawConn) error { for _, ctl := range Controllers { if err := ctl(network, address, c); err != nil { if goerrors.Is(err, ErrControllerRefused) { return err } errors.LogInfoInner(ctx, err, "failed to apply external controller") } } return c.Control(func(fd uintptr) { if sockopt != nil { if err := applyOutboundSocketOptions(network, destAddr.String(), fd, sockopt); err != nil { errors.LogInfo(ctx, err, "failed to apply socket options") } } }) } packetConn, err := lc.ListenPacket(ctx, srcAddr.Network(), srcAddr.String()) if err != nil { return nil, err } return &PacketConnWrapper{ PacketConn: packetConn, Dest: destAddr, }, nil } // Chrome defaults keepAliveConfig := net.KeepAliveConfig{ Enable: true, Idle: 45 * time.Second, Interval: 45 * time.Second, Count: -1, } keepAlive := time.Duration(0) if sockopt != nil { if sockopt.TcpKeepAliveIdle*sockopt.TcpKeepAliveInterval < 0 { return nil, errors.New("invalid TcpKeepAliveIdle or TcpKeepAliveInterval value: ", sockopt.TcpKeepAliveIdle, " ", sockopt.TcpKeepAliveInterval) } if sockopt.TcpKeepAliveIdle < 0 || sockopt.TcpKeepAliveInterval < 0 { keepAlive = -1 keepAliveConfig.Enable = false } if sockopt.TcpKeepAliveIdle > 0 { keepAliveConfig.Idle = time.Duration(sockopt.TcpKeepAliveIdle) * time.Second } if sockopt.TcpKeepAliveInterval > 0 { keepAliveConfig.Interval = time.Duration(sockopt.TcpKeepAliveInterval) * time.Second } } dialer := &net.Dialer{ Timeout: time.Second * 16, LocalAddr: resolveSrcAddr(dest.Network, src), KeepAlive: keepAlive, KeepAliveConfig: keepAliveConfig, } if sockopt != nil || len(Controllers) > 0 { if sockopt != nil && sockopt.TcpMptcp { dialer.SetMultipathTCP(true) } dialer.Control = func(network, address string, c syscall.RawConn) error { for _, ctl := range Controllers { if err := ctl(network, address, c); err != nil { if goerrors.Is(err, ErrControllerRefused) { return err } errors.LogInfoInner(ctx, err, "failed to apply external controller") } } return c.Control(func(fd uintptr) { if sockopt != nil { if err := applyOutboundSocketOptions(network, address, fd, sockopt); err != nil { errors.LogInfoInner(ctx, err, "failed to apply socket options") } } }) } } return dialer.DialContext(ctx, dest.Network.SystemString(), dest.NetAddr()) } func (d *DefaultSystemDialer) DestIpAddress() net.IP { return nil } type PacketConnWrapper struct { net.PacketConn Dest net.Addr } func (c *PacketConnWrapper) Read(p []byte) (int, error) { n, _, err := c.PacketConn.ReadFrom(p) return n, err } func (c *PacketConnWrapper) Write(p []byte) (int, error) { return c.PacketConn.WriteTo(p, c.Dest) } func (c *PacketConnWrapper) RemoteAddr() net.Addr { return c.Dest } type SystemDialerAdapter interface { Dial(network string, address string) (net.Conn, error) } type SimpleSystemDialer struct { adapter SystemDialerAdapter } func WithAdapter(dialer SystemDialerAdapter) SystemDialer { return &SimpleSystemDialer{ adapter: dialer, } } func (v *SimpleSystemDialer) Dial(ctx context.Context, src net.Address, dest net.Destination, sockopt *SocketConfig) (net.Conn, error) { return v.adapter.Dial(dest.Network.SystemString(), dest.NetAddr()) } func (d *SimpleSystemDialer) DestIpAddress() net.IP { return nil } // UseAlternativeSystemDialer replaces the current system dialer with a given one. // Caller must ensure there is no race condition. // // xray:api:stable func UseAlternativeSystemDialer(dialer SystemDialer) { if dialer == nil { dialer = &DefaultSystemDialer{} } effectiveSystemDialer = dialer } // RegisterDialerController adds a controller to the effective system dialer. // The controller can be used to operate on file descriptors before they are put into use. // It only works when effective dialer is the default dialer. // // xray:api:beta func RegisterDialerController(ctl func(network, address string, c syscall.RawConn) error) error { if ctl == nil { return errors.New("nil listener controller") } ControllersLock.Lock() Controllers = append(Controllers, ctl) ControllersLock.Unlock() _, ok := effectiveSystemDialer.(*DefaultSystemDialer) if !ok { return errors.New("RegisterListenerController not supported in custom dialer") } return nil } type FakePacketConn struct { net.Conn } func (c *FakePacketConn) ReadFrom(p []byte) (n int, addr net.Addr, err error) { n, err = c.Read(p) return n, &net.UDPAddr{IP: c.Conn.RemoteAddr().(*net.TCPAddr).IP, Port: c.Conn.RemoteAddr().(*net.TCPAddr).Port}, err } func (c *FakePacketConn) WriteTo(p []byte, _ net.Addr) (n int, err error) { return c.Write(p) } func (c *FakePacketConn) LocalAddr() net.Addr { return &net.UDPAddr{IP: c.Conn.LocalAddr().(*net.TCPAddr).IP, Port: c.Conn.LocalAddr().(*net.TCPAddr).Port} } ============================================================================== transport/internet/spanghew_controller_refused_test.go (изменённый файл целиком) package internet import ( "context" goerrors "errors" "syscall" "testing" "github.com/xtls/xray-core/common/net" ) // Spanghew: a controller that returns ErrControllerRefused fails the dial - TCP and UDP; any other // error of a controller is logged and the dial goes on, as upstream. func TestSpanghewControllerRefusesDial(t *testing.T) { listener, err := net.Listen("tcp", "127.0.0.1:0") if err != nil { t.Fatal(err) } defer listener.Close() go func() { for { conn, err := listener.Accept() if err != nil { return } conn.Close() } }() port := net.Port(listener.Addr().(*net.TCPAddr).Port) destinations := []net.Destination{ net.TCPDestination(net.LocalHostIP, port), net.UDPDestination(net.LocalHostIP, port), } ControllersLock.Lock() saved := Controllers ControllersLock.Unlock() defer func() { ControllersLock.Lock() Controllers = saved ControllersLock.Unlock() }() set := func(result error) { ControllersLock.Lock() Controllers = []func(network, address string, c syscall.RawConn) error{ func(string, string, syscall.RawConn) error { return result }, } ControllersLock.Unlock() } dialer := &DefaultSystemDialer{} set(ErrControllerRefused) for _, dest := range destinations { conn, err := dialer.Dial(context.Background(), nil, dest, nil) if conn != nil { conn.Close() } if !goerrors.Is(err, ErrControllerRefused) { t.Errorf("%s: refused controller, dial error: %v", dest, err) } } set(goerrors.New("could not prepare the socket")) for _, dest := range destinations { conn, err := dialer.Dial(context.Background(), nil, dest, nil) if err != nil { t.Errorf("%s: another controller error must not stop the dial: %v", dest, err) continue } conn.Close() } } ============================================================================== common/net/find_process_windows.go (изменённый файл целиком) //go:build windows package net import ( "net" "net/netip" "path/filepath" "strings" "sync" "syscall" "unsafe" "golang.org/x/sys/windows" "github.com/xtls/xray-core/common/errors" ) const ( tcpTableFunc = "GetExtendedTcpTable" tcpTablePidConn = 4 udpTableFunc = "GetExtendedUdpTable" udpTablePid = 1 ) var ( getExTCPTable uintptr getExUDPTable uintptr once sync.Once initErr error ) func initWin32API() error { h, err := windows.LoadLibrary("iphlpapi.dll") if err != nil { return errors.New("LoadLibrary iphlpapi.dll failed").Base(err) } getExTCPTable, err = windows.GetProcAddress(h, tcpTableFunc) if err != nil { return errors.New("GetProcAddress of ", tcpTableFunc, " failed").Base(err) } getExUDPTable, err = windows.GetProcAddress(h, udpTableFunc) if err != nil { return errors.New("GetProcAddress of ", udpTableFunc, " failed").Base(err) } return nil } func FindProcess(network, srcIP string, srcPort uint16, destIP string, destPort uint16) (PID int, Name string, AbsolutePath string, err error) { once.Do(func() { initErr = initWin32API() }) if initErr != nil { return 0, "", "", initErr } isLocal, err := IsLocal(net.ParseIP(srcIP)) if err != nil { return 0, "", "", errors.New("failed to determine if address is local: ", err) } if !isLocal { return 0, "", "", ErrNotLocal } if network != "tcp" && network != "udp" { panic("Unsupported network type for process lookup.") } var class int var fn uintptr switch network { case "tcp": fn = getExTCPTable class = tcpTablePidConn case "udp": fn = getExUDPTable class = udpTablePid default: panic("Unsupported network type for process lookup.") } ip := net.ParseIP(srcIP) port := int(srcPort) addr, ok := netip.AddrFromSlice(ip) if !ok { return 0, "", "", errors.New("invalid IP address") } addr = addr.Unmap() family := windows.AF_INET if addr.Is6() { family = windows.AF_INET6 } buf, err := getTransportTable(fn, family, class) if err != nil { return 0, "", "", err } networkType := Network_TCP if network == "udp" { networkType = Network_UDP } familyType := AddressFamilyIPv4 if addr.Is6() { familyType = AddressFamilyIPv6 } s := newSearcher(networkType, familyType) pid, err := s.Search(buf, addr, uint16(port)) if err != nil { return 0, "", "", err } NameWithPath, err := getExecPathFromPID(pid) // Spanghew: у программы, запущенной по короткому пути 8.3 (C:\PROGRA~1\...), Windows отдаёт путь с короткими // именами, а правила process — по длинным: без этого выбранная программа шла бы мимо VPN. if err == nil && strings.Contains(NameWithPath, "~") { NameWithPath = longPath(NameWithPath) } NameWithPath = filepath.ToSlash(NameWithPath) // drop .exe and path nameSplit := strings.Split(NameWithPath, "/") procName := nameSplit[len(nameSplit)-1] procName = strings.TrimSuffix(procName, ".exe") return int(pid), procName, NameWithPath, err } type searcher struct { itemSize int port int ip int ipSize int pid int tcpState int } func (s *searcher) Search(b []byte, ip netip.Addr, port uint16) (uint32, error) { n := int(readNativeUint32(b[:4])) itemSize := s.itemSize for i := range n { row := b[4+itemSize*i : 4+itemSize*(i+1)] if s.tcpState >= 0 { tcpState := readNativeUint32(row[s.tcpState : s.tcpState+4]) // MIB_TCP_STATE_ESTAB, only check established connections for TCP if tcpState != 5 { continue } } // according to MSDN, only the lower 16 bits of dwLocalPort are used and the port number is in network endian. // this field can be illustrated as follows depends on different machine endianess: // little endian: [ MSB LSB 0 0 ] interpret as native uint32 is ((LSB<<8)|MSB) // big endian: [ 0 0 MSB LSB ] interpret as native uint32 is ((MSB<<8)|LSB) // so we need an syscall.Ntohs on the lower 16 bits after read the port as native uint32 srcPort := syscall.Ntohs(uint16(readNativeUint32(row[s.port : s.port+4]))) if srcPort != port { continue } srcIP, _ := netip.AddrFromSlice(row[s.ip : s.ip+s.ipSize]) srcIP = srcIP.Unmap() // windows binds an unbound udp socket to 0.0.0.0/[::] while first sendto if ip != srcIP && (!srcIP.IsUnspecified() || s.tcpState != -1) { continue } pid := readNativeUint32(row[s.pid : s.pid+4]) return pid, nil } return 0, errors.New("not found") } func newSearcher(network Network, family AddressFamily) *searcher { var itemSize, port, ip, ipSize, pid int tcpState := -1 switch network { case Network_TCP: if family == AddressFamilyIPv4 { // struct MIB_TCPROW_OWNER_PID itemSize, port, ip, ipSize, pid, tcpState = 24, 8, 4, 4, 20, 0 } if family == AddressFamilyIPv6 { // struct MIB_TCP6ROW_OWNER_PID itemSize, port, ip, ipSize, pid, tcpState = 56, 20, 0, 16, 52, 48 } case Network_UDP: if family == AddressFamilyIPv4 { // struct MIB_UDPROW_OWNER_PID itemSize, port, ip, ipSize, pid = 12, 4, 0, 4, 8 } if family == AddressFamilyIPv6 { // struct MIB_UDP6ROW_OWNER_PID itemSize, port, ip, ipSize, pid = 28, 20, 0, 16, 24 } } return &searcher{ itemSize: itemSize, port: port, ip: ip, ipSize: ipSize, pid: pid, tcpState: tcpState, } } func getTransportTable(fn uintptr, family int, class int) ([]byte, error) { for size, buf := uint32(8), make([]byte, 8); ; { ptr := unsafe.Pointer(&buf[0]) err, _, _ := syscall.Syscall6(fn, 6, uintptr(ptr), uintptr(unsafe.Pointer(&size)), 0, uintptr(family), uintptr(class), 0) switch err { case 0: return buf, nil case uintptr(syscall.ERROR_INSUFFICIENT_BUFFER): buf = make([]byte, size) default: return nil, errors.New("syscall error: ", int(err)) } } } func readNativeUint32(b []byte) uint32 { return *(*uint32)(unsafe.Pointer(&b[0])) } // Spanghew: длинная форма пути (GetLongPathNameW); не вышло — путь как есть. func longPath(path string) string { p, err := windows.UTF16PtrFromString(path) if err != nil { return path } buf := make([]uint16, syscall.MAX_LONG_PATH) n, err := windows.GetLongPathName(p, &buf[0], uint32(len(buf))) if err != nil || n == 0 || n >= uint32(len(buf)) { return path } return windows.UTF16ToString(buf[:n]) } func getExecPathFromPID(pid uint32) (string, error) { // kernel process starts with a colon in order to distinguish with normal processes switch pid { case 0: // reserved pid for system idle process return ":System Idle Process", nil case 4: // reserved pid for windows kernel image return ":System", nil } h, err := windows.OpenProcess(windows.PROCESS_QUERY_LIMITED_INFORMATION, false, pid) if err != nil { return "", err } defer windows.CloseHandle(h) buf := make([]uint16, syscall.MAX_LONG_PATH) size := uint32(len(buf)) err = windows.QueryFullProcessImageName(h, 0, &buf[0], &size) if err != nil { return "", err } return syscall.UTF16ToString(buf[:size]), nil } ============================================================================== proxy/tun/tun_darwin.go (изменённый файл целиком) //go:build darwin package tun import ( "context" "errors" "fmt" "net" "net/netip" "os" "strconv" "sync" "sync/atomic" "time" "unsafe" "github.com/xtls/xray-core/common/buf" xerrors "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/platform" "golang.org/x/net/route" "golang.org/x/sys/unix" "gvisor.dev/gvisor/pkg/buffer" "gvisor.dev/gvisor/pkg/tcpip" "gvisor.dev/gvisor/pkg/tcpip/stack" ) const ( utunControlName = "com.apple.net.utun_control" sysprotoControl = 2 defaultDarwinGateway = "169.254.10.1/30" utunHeaderSize = 4 UTUN_OPT_IFNAME = 2 ) const ( SIOCAIFADDR6 = 2155899162 // netinet6/in6_var.h IN6_IFF_NODAD = 0x0020 // netinet6/in6_var.h IN6_IFF_SECURED = 0x0400 // netinet6/in6_var.h ND6_INFINITE_LIFETIME = 0xFFFFFFFF // netinet6/nd6.h ) type DarwinTun struct { tunFile *os.File options *Config tunFd int ownsFd bool // true for macOS (we created the fd), false for iOS (fd from system) // Genuinely blocks Wait() until tunFd is readable, instead of the // previous procyield-only busy-spin (dispatchLoop in // stack_gvisor_endpoint.go calls ReadPacket() then Wait() in a tight // loop with no other throttling whenever the queue is empty -- with // only procyield(1), that pins a full CPU core for as long as the // tunnel is up, observed causing severe device heating/thermal // shutdown). nil if kqueue setup failed, in which case Wait() falls // back to a bounded time.Sleep instead. See waitKqueue's own doc // comment for why this is a dedicated type rather than a bare fd. waitKq *waitKqueue routeMonitor *os.File routeMonitorOnce sync.Once systemRoutes []netip.Prefix gateway netip.Prefix } // waitKqueue owns a kqueue fd used by DarwinTun.Wait() to block on // read-readiness. Closing and waiting can race from different goroutines // (Close() from the caller that tears down the tunnel, Wait() from // dispatchLoop's own goroutine) -- reviewer feedback on XTLS/Xray-core#6580 // found that a bare `int` fd field let Close() race Wait()'s use of the // same fd number, and on Darwin a closed fd number can be reused by an // unrelated concurrent open() before Wait() gets to call Kevent on it, // so Wait() could end up polling (or Close() could end up closing) a // completely unrelated file descriptor. This type makes closing // idempotent (sync.Once) and gates every Kevent call behind an atomic // "closed" flag checked immediately before the syscall, so Wait() never // issues a kevent syscall against a fd number that Close() has already // (or is concurrently) invalidated -- there's still a narrow window where // Wait() checks-then-uses the fd, but Close() only actually closes it // after Wait() cannot start a new syscall on it (the flag is set first, // synchronized with acquire/release semantics), which is sufficient since // Wait()'s Kevent call itself is what's being raced, not a fd read/write. type waitKqueue struct { fd int closed atomic.Bool once sync.Once } // newWaitKqueue creates a kqueue registered for read-readiness on fd, for // Wait() to block on. Returns nil if anything fails, so callers can fall // back to a bounded sleep rather than error out of NewTun over what is // purely a CPU-efficiency concern. func newWaitKqueue(fd int) *waitKqueue { kq, err := unix.Kqueue() if err != nil { return nil } _, err = unix.Kevent(kq, []unix.Kevent_t{{ Ident: uint64(fd), Filter: unix.EVFILT_READ, Flags: unix.EV_ADD | unix.EV_ENABLE, }}, nil, nil) if err != nil { _ = unix.Close(kq) return nil } return &waitKqueue{fd: kq} } // wait blocks until the registered fd is readable, timeout elapses, or a // benign interrupt occurs -- all three are "this kqueue is still healthy, // the caller should just try again" and return true; the caller // (DarwinTun.Wait) doesn't need to distinguish them since it always calls // ReadPacket() right after anyway, and that already handles "nothing was // actually there" via ErrQueueEmpty. Returns false only when the kqueue // itself is no longer usable -- already closed, or the kevent syscall // failed for a reason other than EINTR -- see its own call site in // DarwinTun.Wait for why a persistent failure must not be silently // retried forever (reviewer feedback, XTLS/Xray-core#6580 P2). func (w *waitKqueue) wait(timeout time.Duration) (ok bool) { if w.closed.Load() { return false } events := make([]unix.Kevent_t, 1) ts := unix.NsecToTimespec(timeout.Nanoseconds()) _, err := unix.Kevent(w.fd, nil, events, &ts) if err != nil { return errors.Is(err, unix.EINTR) } return true } // close marks the kqueue as unusable (so any Wait() call that hasn't yet // entered the kevent syscall bails out instead) and closes the underlying // fd exactly once, regardless of how many times close is called or // whether it races a Wait() already inside its kevent syscall (that call // either completes against the still-open fd or returns an error safely // -- either way, no other goroutine can be handed this fd number in // between the atomic flag flip and the actual close, since nothing else // in this type ever creates a new kqueue with the same field). func (w *waitKqueue) close() { w.once.Do(func() { w.closed.Store(true) _ = unix.Close(w.fd) }) } var ( _ Tun = (*DarwinTun)(nil) _ GVisorDevice = (*DarwinTun)(nil) ) func NewTun(options *Config) (Tun, error) { // Check if fd is provided via environment (iOS mode) fdStr := platform.NewEnvFlag(platform.TunFdKey).GetValue(func() string { return "" }) if fdStr != "" { // iOS: use provided fd from NetworkExtension fd, err := strconv.Atoi(fdStr) if err != nil { return nil, err } if err = unix.SetNonblock(fd, true); err != nil { return nil, err } // Spanghew: work on our own copy of the NetworkExtension fd. os.File closes its fd from a GC finalizer; // after a server change (core stop + start on the same utun) the old core's file would close the fd // number the new core reads. The copy is ours alone: Close() closes it and wakes the old reader. own, err := unix.Dup(fd) if err != nil { return nil, err } unix.CloseOnExec(own) return &DarwinTun{ tunFile: os.NewFile(uintptr(own), "utun"), options: options, tunFd: own, ownsFd: false, waitKq: newWaitKqueue(own), }, nil } // macOS: create our own utun interface tunFile, err := open(options.Name) if err != nil { return nil, err } gateway, err := selectDarwinGateway(options.Gateway) if err != nil { _ = tunFile.Close() return nil, err } err = setup(options.Name, options.MTU, gateway) if err != nil { _ = tunFile.Close() return nil, err } return &DarwinTun{ tunFile: tunFile, options: options, tunFd: int(tunFile.Fd()), ownsFd: true, waitKq: newWaitKqueue(int(tunFile.Fd())), gateway: gateway, }, nil } func (t *DarwinTun) Start() error { if !t.ownsFd { return nil } if err := t.setSystemRoutes(); err != nil { return err } if updater != nil { fd, err := unix.Socket(unix.AF_ROUTE, unix.SOCK_RAW, 0) if err != nil { _ = t.unsetSystemRoutes() return err } t.routeMonitor = os.NewFile(uintptr(fd), "xray-route-monitor") go t.monitorRouteChanges() } return nil } func (t *DarwinTun) Close() error { t.routeMonitorOnce.Do(func() { if t.routeMonitor != nil { _ = t.routeMonitor.Close() } }) if t.waitKq != nil { t.waitKq.close() } routeErr := t.unsetSystemRoutes() // Spanghew: the fd from NetworkExtension is a copy made in NewTun — close it too; the original stays with // the system. return xerrors.Combine(routeErr, t.tunFile.Close()) } func (t *DarwinTun) monitorRouteChanges() { buffer := make([]byte, 64*1024) for { if _, err := t.routeMonitor.Read(buffer); err != nil { if !errors.Is(err, os.ErrClosed) { xerrors.LogInfoInner(context.Background(), err, "[tun] failed to monitor route changes") } return } if updater != nil { updater.Update() } } } func (t *DarwinTun) Name() (string, error) { return unix.GetsockoptString(t.tunFd, sysprotoControl, UTUN_OPT_IFNAME) } func (t *DarwinTun) Index() (int, error) { name, err := t.Name() if err != nil { return 0, err } iface, err := net.InterfaceByName(name) if err != nil { return 0, err } return iface.Index, nil } // WritePacket implements GVisorDevice method to write one packet to the tun device func (t *DarwinTun) WritePacket(packet *stack.PacketBuffer) tcpip.Error { // request memory to write from reusable buffer pool b := buf.NewWithSize(int32(t.options.MTU) + utunHeaderSize) defer b.Release() // prepare Darwin specific packet header _, _ = b.Write([]byte{0x0, 0x0, 0x0, 0x0}) // copy the bytes of slices that compose the packet into the allocated buffer for _, packetElement := range packet.AsSlices() { _, _ = b.Write(packetElement) } // fill Darwin specific header from the first raw packet byte, that we can access now var family byte switch b.Byte(4) >> 4 { case 4: family = unix.AF_INET case 6: family = unix.AF_INET6 default: return &tcpip.ErrAborted{} } b.SetByte(3, family) if _, err := t.tunFile.Write(b.Bytes()); err != nil { if errors.Is(err, unix.EAGAIN) { return &tcpip.ErrWouldBlock{} } return &tcpip.ErrAborted{} } return nil } // ReadPacket implements GVisorDevice method to read one packet from the tun device // It is expected that the method will not block, rather return ErrQueueEmpty when there is nothing on the line, // which will make the stack call Wait which should implement desired push-back func (t *DarwinTun) ReadPacket() (byte, *stack.PacketBuffer, error) { // request memory to write from reusable buffer pool b := buf.NewWithSize(int32(t.options.MTU) + utunHeaderSize) // read the bytes to the interface file n, err := b.ReadFrom(t.tunFile) if errors.Is(err, unix.EAGAIN) || errors.Is(err, unix.EINTR) { b.Release() return 0, nil, ErrQueueEmpty } if err != nil { b.Release() return 0, nil, err } // discard empty or sub-empty packets if n <= utunHeaderSize { b.Release() return 0, nil, ErrQueueEmpty } // network protocol version from first byte of the raw packet, the one that follows Darwin specific header version := b.Byte(utunHeaderSize) >> 4 packetBuffer := buffer.MakeWithData(b.BytesFrom(utunHeaderSize)) return version, stack.NewPacketBuffer(stack.PacketBufferOptions{ Payload: packetBuffer, IsForwardedPacket: true, OnRelease: func() { b.Release() }, }), nil } // Wait blocks until tunFd is readable (or a short timeout elapses), rather // than spinning the CPU -- see the waitKq field's own doc comment. A bounded // timeout (not an indefinite wait) keeps this responsive to a Close() that // happens to race a call already parked here. // // Reviewer feedback (XTLS/Xray-core#6580, P2): the original version // discarded every error from the underlying kevent syscall. dispatchLoop // (stack_gvisor_endpoint.go) calls ReadPacket() then Wait() in an // unconditional tight loop -- if kevent started failing at runtime for a // persistent reason (not just a benign EINTR), Wait() returning // immediately every time reintroduces exactly the busy-spin this whole // change exists to remove, just routed through a failing syscall instead // of procyield. waitKq.wait's own bool return distinguishes "genuinely // interrupted, try again" from "this kqueue is unusable now" -- Wait() // permanently falls back to the sleep path once that happens, rather than // retrying the same broken kqueue forever. func (t *DarwinTun) Wait() { if t.waitKq != nil && t.waitKq.wait(time.Second) { return } if t.waitKq != nil { // Persistent kevent failure (not a benign EINTR, and not just // "the 1s timeout elapsed with nothing to read" -- wait() already // returned true for both of those cases above). Stop trusting // this kqueue for the rest of this DarwinTun's lifetime instead of // re-attempting a syscall that's already shown it won't succeed. t.waitKq.close() t.waitKq = nil } // Reviewer feedback (XTLS/Xray-core#6580): procyield here is the same // busy-spin this whole change exists to remove, just gated behind an // edge case (kqueue setup failing, which practically never happens on // real Darwin systems, or having just failed permanently above) // instead of always -- a genuine bounded sleep actually yields the CPU // instead of being a near-instant scheduler hint that lets the tight // dispatchLoop caller spin just as hot as before. time.Sleep(time.Millisecond) } func (t *DarwinTun) newEndpoint() (stack.LinkEndpoint, error) { return &LinkEndpoint{deviceMTU: t.options.MTU, device: t}, nil } // open the interface, by creating new utunN if in the system and returning its file descriptor func open(name string) (*os.File, error) { ifIndex := -1 _, err := fmt.Sscanf(name, "utun%d", &ifIndex) if err != nil || ifIndex < 0 { return nil, errors.New("interface name must be utunN, where N is a number, e.g. utun9, utun11 and so on") } fd, err := unix.Socket(unix.AF_SYSTEM, unix.SOCK_DGRAM, sysprotoControl) if err != nil { return nil, err } ctlInfo := &unix.CtlInfo{} copy(ctlInfo.Name[:], utunControlName) if err := unix.IoctlCtlInfo(fd, ctlInfo); err != nil { _ = unix.Close(fd) return nil, err } sockaddr := &unix.SockaddrCtl{ ID: ctlInfo.Id, Unit: uint32(ifIndex) + 1, } if err := unix.Connect(fd, sockaddr); err != nil { _ = unix.Close(fd) return nil, err } if err := unix.SetNonblock(fd, true); err != nil { _ = unix.Close(fd) return nil, err } return os.NewFile(uintptr(fd), name), nil } // setup the interface by name func setup(name string, MTU uint32, gateway netip.Prefix) error { if err := setMTU(name, MTU); err != nil { return err } /* * Darwin routing require tunnel type interface to have local and remote address, to be routable. * To simplify inevitable task, assign the interface static ip address. */ if err := setIPAddress(name, gateway); err != nil { return err } return nil } func selectDarwinGateway(configured []string) (netip.Prefix, error) { if len(configured) == 0 { return netip.ParsePrefix(defaultDarwinGateway) } for _, value := range configured { prefix, err := netip.ParsePrefix(value) if err != nil { return netip.Prefix{}, xerrors.New("invalid macOS gateway ", value).Base(err) } if !prefix.Addr().Is4() { continue } local, ok := nextDarwinLocalIPv4(prefix) if !ok || !prefix.Contains(local) { return netip.Prefix{}, xerrors.New("macOS gateway ", value, " must contain at least one usable local IPv4 address after the gateway address") } return prefix, nil } return netip.Prefix{}, xerrors.New("macOS gateway requires at least one IPv4 prefix") } func nextDarwinLocalIPv4(gateway netip.Prefix) (netip.Addr, bool) { local4 := gateway.Addr().As4() for i := len(local4) - 1; i >= 0; i-- { local4[i]++ if local4[i] != 0 { return netip.AddrFrom4(local4), true } } return netip.Addr{}, false } // setMTU sets MTU on the interface by given name func setMTU(name string, mtu uint32) error { socket, err := unix.Socket(unix.AF_INET, unix.SOCK_DGRAM, 0) if err != nil { return err } defer unix.Close(socket) ifr := unix.IfreqMTU{MTU: int32(mtu)} copy(ifr.Name[:], name) return unix.IoctlSetIfreqMTU(socket, &ifr) } type ifAliasReq4 struct { Name [unix.IFNAMSIZ]byte Addr unix.RawSockaddrInet4 Dstaddr unix.RawSockaddrInet4 Mask unix.RawSockaddrInet4 } type ifAliasReq6 struct { Name [unix.IFNAMSIZ]byte Addr unix.RawSockaddrInet6 Dstaddr unix.RawSockaddrInet6 Mask unix.RawSockaddrInet6 Flags uint32 Lifetime addrLifetime6 } type addrLifetime6 struct { Expire float64 Preferred float64 Vltime uint32 Pltime uint32 } // setIPAddress sets ipv4 and ipv6 addresses to the interface, required for the routing to work func setIPAddress(name string, gateway netip.Prefix) error { socket4, err := unix.Socket(unix.AF_INET, unix.SOCK_DGRAM, 0) if err != nil { return err } defer unix.Close(socket4) // assume local ip address is next one from the remote address local, ok := nextDarwinLocalIPv4(gateway) if !ok || !gateway.Contains(local) { return xerrors.New("macOS gateway ", gateway.String(), " must contain at least one usable local IPv4 address after the gateway address") } local4 := local.As4() // fill the configuration for ipv4 ifReq4 := ifAliasReq4{ Addr: unix.RawSockaddrInet4{ Len: unix.SizeofSockaddrInet4, Family: unix.AF_INET, Addr: local4, }, Dstaddr: unix.RawSockaddrInet4{ Len: unix.SizeofSockaddrInet4, Family: unix.AF_INET, Addr: gateway.Addr().As4(), }, Mask: unix.RawSockaddrInet4{ Len: unix.SizeofSockaddrInet4, Family: unix.AF_INET, Addr: netip.MustParseAddr(net.IP(net.CIDRMask(gateway.Bits(), 32)).String()).As4(), }, } copy(ifReq4.Name[:], name) if err = ioctlPtr(socket4, unix.SIOCAIFADDR, unsafe.Pointer(&ifReq4)); err != nil { return os.NewSyscallError("SIOCAIFADDR", err) } socket6, err := unix.Socket(unix.AF_INET6, unix.SOCK_DGRAM, 0) if err != nil { return err } defer unix.Close(socket6) // link-local ipv6 address with suffix from ipv6 local6 := netip.AddrFrom16([16]byte{0: 0xfe, 1: 0x80, 12: local4[0], 13: local4[1], 14: local4[2], 15: local4[3]}) // fill the configuration for ipv6 // only link-local address without the destination is enough for it ifReq6 := ifAliasReq6{ Addr: unix.RawSockaddrInet6{ Len: unix.SizeofSockaddrInet6, Family: unix.AF_INET6, Addr: local6.As16(), }, Mask: unix.RawSockaddrInet6{ Len: unix.SizeofSockaddrInet6, Family: unix.AF_INET6, Addr: netip.MustParseAddr(net.IP(net.CIDRMask(64, 128)).String()).As16(), }, Flags: IN6_IFF_NODAD, Lifetime: addrLifetime6{ Vltime: ND6_INFINITE_LIFETIME, Pltime: ND6_INFINITE_LIFETIME, }, } // assign link-local ipv6 address to the interface. // this will additionally trigger OS level autoconfiguration, which will result two different link-local // addresses - the requested one, and autoconfigured one. // this really has no known side effects, just look excessive. and actually considered pretty normal way to // enable the ipv6 on the interface by macOS concepts. copy(ifReq6.Name[:], name) if err = ioctlPtr(socket6, SIOCAIFADDR6, unsafe.Pointer(&ifReq6)); err != nil { return os.NewSyscallError("SIOCAIFADDR6", err) } return nil } func ioctlPtr(fd int, req uint, arg unsafe.Pointer) error { _, _, errno := unix.Syscall(unix.SYS_IOCTL, uintptr(fd), uintptr(req), uintptr(arg)) if errno != 0 { return errno } return nil } func setinterface(network, address string, fd uintptr, iface *net.Interface) error { var err1, err2 error switch network { case "tcp6", "udp6", "ip6": err1 = unix.SetsockoptInt(int(fd), unix.IPPROTO_IPV6, unix.IPV6_BOUND_IF, iface.Index) fallthrough case "tcp4", "udp4", "ip4": err2 = unix.SetsockoptInt(int(fd), unix.IPPROTO_IP, unix.IP_BOUND_IF, iface.Index) default: panic(network + " " + address) } return errors.Join(err1, err2) } func findOutboundInterface(tunIndex int, fixedName string) (*net.Interface, error) { if fixedName != "" { iface, err := net.InterfaceByName(fixedName) if err != nil { return nil, err } if iface.Index == tunIndex { return nil, errors.New("outbound interface cannot be the TUN interface") } return iface, nil } rib, err := route.FetchRIB(unix.AF_UNSPEC, route.RIBTypeRoute, 0) if err != nil { return nil, err } messages, err := route.ParseRIB(route.RIBTypeRoute, rib) if err != nil { return nil, err } var ipv6Index int for _, message := range messages { routeMessage, ok := message.(*route.RouteMessage) if !ok || routeMessage.Index == tunIndex { continue } if routeMessage.Flags&unix.RTF_UP == 0 || routeMessage.Flags&unix.RTF_GATEWAY == 0 { continue } family, ok := defaultRouteFamily(routeMessage) if !ok { continue } if family == unix.AF_INET { return usableDarwinInterface(routeMessage.Index) } if family == unix.AF_INET6 && ipv6Index == 0 { ipv6Index = routeMessage.Index } } if ipv6Index != 0 { return usableDarwinInterface(ipv6Index) } return nil, errors.New("default route not found") } func defaultRouteFamily(message *route.RouteMessage) (int, bool) { if len(message.Addrs) <= unix.RTAX_NETMASK { return 0, false } switch destination := message.Addrs[unix.RTAX_DST].(type) { case *route.Inet4Addr: mask, ok := message.Addrs[unix.RTAX_NETMASK].(*route.Inet4Addr) if !ok || destination.IP != netip.IPv4Unspecified().As4() { return 0, false } ones, bits := net.IPMask(mask.IP[:]).Size() return unix.AF_INET, ones == 0 && bits == 32 case *route.Inet6Addr: mask, ok := message.Addrs[unix.RTAX_NETMASK].(*route.Inet6Addr) if !ok || destination.IP != netip.IPv6Unspecified().As16() { return 0, false } ones, bits := net.IPMask(mask.IP[:]).Size() return unix.AF_INET6, ones == 0 && bits == 128 default: return 0, false } } func usableDarwinInterface(index int) (*net.Interface, error) { iface, err := net.InterfaceByIndex(index) if err != nil { return nil, err } if iface.Flags&net.FlagUp == 0 || iface.Flags&net.FlagLoopback != 0 { return nil, errors.New("default route interface is not usable") } return iface, nil } func (t *DarwinTun) setSystemRoutes() error { routes, err := buildDarwinSystemRoutes(t.options.AutoSystemRoutingTable) if err != nil { return err } if len(routes) == 0 { return nil } tunIndex, err := t.Index() if err != nil { return err } for _, destination := range routes { if err := execDarwinRoute(unix.RTM_ADD, tunIndex, destination, t.gateway); err != nil { _ = t.unsetSystemRoutes() return xerrors.New("failed to add system route ", destination).Base(err) } t.systemRoutes = append(t.systemRoutes, destination) } return nil } func (t *DarwinTun) unsetSystemRoutes() error { var errs []error tunIndex, indexErr := t.Index() if indexErr != nil && len(t.systemRoutes) > 0 { errs = append(errs, indexErr) } for i := len(t.systemRoutes) - 1; i >= 0; i-- { destination := t.systemRoutes[i] if err := execDarwinRoute(unix.RTM_DELETE, tunIndex, destination, t.gateway); err != nil && !errors.Is(err, unix.ESRCH) { errs = append(errs, xerrors.New("failed to delete system route ", destination).Base(err)) } } t.systemRoutes = nil return xerrors.Combine(errs...) } func buildDarwinSystemRoutes(configured []string) ([]netip.Prefix, error) { routes := make([]netip.Prefix, 0, len(configured)) seen := make(map[netip.Prefix]struct{}) appendRoute := func(prefix netip.Prefix) { prefix = prefix.Masked() if _, found := seen[prefix]; found { return } seen[prefix] = struct{}{} routes = append(routes, prefix) } for _, value := range configured { prefix, err := netip.ParsePrefix(value) if err != nil { return nil, xerrors.New("invalid system route ", value).Base(err) } prefix = prefix.Masked() if prefix.Bits() == 0 { for _, protected := range darwinProtectedDefaultRoutes(prefix.Addr().Is4()) { appendRoute(protected) } continue } appendRoute(prefix) } return routes, nil } func darwinProtectedDefaultRoutes(ipv4 bool) []netip.Prefix { routes := make([]netip.Prefix, 0, 8) for i := 0; i < 8; i++ { if ipv4 { var address [4]byte address[0] = 1 << i routes = append(routes, netip.PrefixFrom(netip.AddrFrom4(address), 8-i)) } else { var address [16]byte address[0] = 1 << i routes = append(routes, netip.PrefixFrom(netip.AddrFrom16(address), 8-i)) } } return routes } func execDarwinRoute(messageType int, interfaceIndex int, destination netip.Prefix, gateway netip.Prefix) error { message := route.RouteMessage{ Type: messageType, Version: unix.RTM_VERSION, Flags: unix.RTF_STATIC | unix.RTF_GATEWAY, Seq: 1, } if messageType == unix.RTM_ADD { message.Flags |= unix.RTF_UP } if destination.Addr().Is4() { message.Addrs = []route.Addr{ unix.RTAX_DST: &route.Inet4Addr{IP: destination.Addr().As4()}, unix.RTAX_NETMASK: &route.Inet4Addr{IP: prefixMask4(destination.Bits())}, unix.RTAX_GATEWAY: &route.Inet4Addr{IP: gateway.Addr().As4()}, } } else { message.Flags &^= unix.RTF_GATEWAY message.Index = interfaceIndex message.Addrs = []route.Addr{ unix.RTAX_DST: &route.Inet6Addr{IP: destination.Addr().As16()}, unix.RTAX_NETMASK: &route.Inet6Addr{IP: prefixMask6(destination.Bits())}, unix.RTAX_GATEWAY: &route.LinkAddr{Index: interfaceIndex}, } } request, err := message.Marshal() if err != nil { return err } fd, err := unix.Socket(unix.AF_ROUTE, unix.SOCK_RAW, 0) if err != nil { return err } defer unix.Close(fd) _, err = unix.Write(fd, request) return err } func prefixMask4(bits int) [4]byte { var mask [4]byte copy(mask[:], net.CIDRMask(bits, 32)) return mask } func prefixMask6(bits int) [16]byte { var mask [16]byte copy(mask[:], net.CIDRMask(bits, 128)) return mask } ============================================================================== proxy/tun/spanghew_bufsize.go (изменённый файл целиком) package tun import "runtime" // Spanghew: upper limits of the per-connection TCP buffers of the gVisor // stack. Upstream: 8 MiB receive, 6 MiB send, window auto-tuning on. The iOS // Packet Tunnel extension lives under a ~50 MiB memory limit (iPhone 7, // iOS 15), and a few uploads through a slow server could grow the windows past // it - iOS then kills the extension. The TUN side is local (no round trip // time), so smaller windows do not limit the speed. Not below the defaults // (1 MiB both): gVisor refuses a range whose maximum is under its default, // and the stack would not start (1.0.25 on the iPhone, 3.10). var tcpRXBufMaxSize, tcpTXBufMaxSize = tcpBufMaxSizes(runtime.GOOS) func tcpBufMaxSizes(goos string) (rx, tx int) { if goos == "ios" { return 1 << 20, 1 << 20 } return 8 << 20, 6 << 20 } ============================================================================== proxy/tun/spanghew_bufsize_test.go (изменённый файл целиком) package tun import ( "testing" "gvisor.dev/gvisor/pkg/tcpip" "gvisor.dev/gvisor/pkg/tcpip/network/ipv4" "gvisor.dev/gvisor/pkg/tcpip/stack" "gvisor.dev/gvisor/pkg/tcpip/transport/tcp" ) // The ranges must be accepted by the gVisor stack itself, the way Start sets them. func TestSpanghewTCPBufMaxSizes(t *testing.T) { for _, goos := range []string{"ios", "darwin", "android", "windows", "linux"} { rx, tx := tcpBufMaxSizes(goos) s := stack.New(stack.Options{ NetworkProtocols: []stack.NetworkProtocolFactory{ipv4.NewProtocol}, TransportProtocols: []stack.TransportProtocolFactory{tcp.NewProtocol}, }) rxOpt := tcpip.TCPReceiveBufferSizeRangeOption{Min: tcpRXBufMinSize, Default: tcpRXBufDefSize, Max: rx} if err := s.SetTransportProtocolOption(tcp.ProtocolNumber, &rxOpt); err != nil { t.Errorf("%s: receive range refused: %v", goos, err) } txOpt := tcpip.TCPSendBufferSizeRangeOption{Min: tcpTXBufMinSize, Default: tcpTXBufDefSize, Max: tx} if err := s.SetTransportProtocolOption(tcp.ProtocolNumber, &txOpt); err != nil { t.Errorf("%s: send range refused: %v", goos, err) } s.Close() } if rx, tx := tcpBufMaxSizes("ios"); rx != 1<<20 || tx != 1<<20 { t.Fatalf("ios: %d, %d", rx, tx) } if rx, tx := tcpBufMaxSizes("darwin"); rx != 8<<20 || tx != 6<<20 { t.Fatalf("darwin: %d, %d", rx, tx) } } ============================================================================== proxy/tun/stack_gvisor_endpoint.go (изменённый файл целиком) package tun import ( "context" "errors" xerrors "github.com/xtls/xray-core/common/errors" "gvisor.dev/gvisor/pkg/tcpip" "gvisor.dev/gvisor/pkg/tcpip/header" "gvisor.dev/gvisor/pkg/tcpip/stack" ) var ErrQueueEmpty = errors.New("queue is empty") type GVisorDevice interface { WritePacket(packet *stack.PacketBuffer) tcpip.Error ReadPacket() (byte, *stack.PacketBuffer, error) Wait() } // LinkEndpoint implements GVisor stack.LinkEndpoint var _ stack.LinkEndpoint = (*LinkEndpoint)(nil) type LinkEndpoint struct { deviceMTU uint32 device GVisorDevice dispatcherCancel context.CancelFunc } func (e *LinkEndpoint) MTU() uint32 { return e.deviceMTU } func (e *LinkEndpoint) SetMTU(_ uint32) { // not Implemented, as it is not expected GVisor will be asking tun device to be modified } func (e *LinkEndpoint) MaxHeaderLength() uint16 { return 0 } func (e *LinkEndpoint) LinkAddress() tcpip.LinkAddress { return "" } func (e *LinkEndpoint) SetLinkAddress(_ tcpip.LinkAddress) { // not Implemented, as it is not expected GVisor will be asking tun device to be modified } func (e *LinkEndpoint) Capabilities() stack.LinkEndpointCapabilities { return stack.CapabilityRXChecksumOffload } func (e *LinkEndpoint) Attach(dispatcher stack.NetworkDispatcher) { if e.dispatcherCancel != nil { e.dispatcherCancel() e.dispatcherCancel = nil } if dispatcher != nil { ctx, cancel := context.WithCancel(context.Background()) go e.dispatchLoop(ctx, dispatcher) e.dispatcherCancel = cancel } } func (e *LinkEndpoint) IsAttached() bool { return e.dispatcherCancel != nil } func (e *LinkEndpoint) Wait() { } func (e *LinkEndpoint) ARPHardwareType() header.ARPHardwareType { return header.ARPHardwareNone } func (e *LinkEndpoint) AddHeader(buffer *stack.PacketBuffer) { // tun interface doesn't have link layer header, it will be added by the OS } func (e *LinkEndpoint) ParseHeader(ptr *stack.PacketBuffer) bool { return true } func (e *LinkEndpoint) Close() { if e.dispatcherCancel != nil { e.dispatcherCancel() e.dispatcherCancel = nil } } func (e *LinkEndpoint) SetOnCloseAction(_ func()) { } func (e *LinkEndpoint) WritePackets(packetBufferList stack.PacketBufferList) (int, tcpip.Error) { var n int var err tcpip.Error for _, packetBuffer := range packetBufferList.AsSlice() { err = e.device.WritePacket(packetBuffer) if err != nil { return n, &tcpip.ErrAborted{} } n++ } return n, nil } func (e *LinkEndpoint) dispatchLoop(ctx context.Context, dispatcher stack.NetworkDispatcher) { var networkProtocolNumber tcpip.NetworkProtocolNumber var version byte var packet *stack.PacketBuffer var err error for { select { case <-ctx.Done(): return default: version, packet, err = e.device.ReadPacket() // on "queue empty", ask device to yield slightly and continue if errors.Is(err, ErrQueueEmpty) { e.device.Wait() continue } // stop dispatcher loop on any other interface failure if err != nil { // Spanghew: the loop used to stop silently — traffic into the TUN then goes nowhere. if ctx.Err() == nil { xerrors.LogWarningInner(ctx, err, "[tun] packet read loop stopped") } e.Attach(nil) return } // extract network protocol number from the packet first byte // (which is returned separately, since it is so incredibly hard to extract one byte from // stack.PacketBuffer without additional memory allocation and full copying it back and forth) switch version { case 4: networkProtocolNumber = header.IPv4ProtocolNumber case 6: networkProtocolNumber = header.IPv6ProtocolNumber default: // discard unknown network protocol packet packet.DecRef() continue } // dispatch the buffer to the stack dispatcher.DeliverNetworkPacket(networkProtocolNumber, packet) // signal the buffer that it can be released packet.DecRef() } } } ============================================================================== xray-freedom-udp-route.patch --- a/proxy/freedom/freedom.go +++ b/proxy/freedom/freedom.go @@ -23,7 +23,10 @@ "github.com/xtls/xray-core/common/task" "github.com/xtls/xray-core/common/utils" "github.com/xtls/xray-core/core" + "github.com/xtls/xray-core/features/outbound" "github.com/xtls/xray-core/features/policy" + "github.com/xtls/xray-core/features/routing" + routing_session "github.com/xtls/xray-core/features/routing/session" "github.com/xtls/xray-core/features/stats" "github.com/xtls/xray-core/proxy" "github.com/xtls/xray-core/transport" @@ -53,6 +56,7 @@ func init() { common.Must(common.RegisterConfig((*Config)(nil), func(ctx context.Context, config interface{}) (interface{}, error) { h := new(Handler) + h.instance = core.FromContext(ctx) if streamSettings, ok := session.StreamSettingsFromContext(ctx).(*internet.MemoryStreamConfig); ok && streamSettings.SocketSettings != nil { h.resolveStrategy = streamSettings.SocketSettings.DomainStrategy h.usesDialerProxy = len(streamSettings.SocketSettings.DialerProxy) > 0 @@ -98,8 +102,89 @@ finalRules []*FinalRule resolveStrategy internet.DomainStrategy usesDialerProxy bool + instance *core.Instance } +// Spanghew: the route of a UDP connection is chosen by its first packet, and freedom then sent every +// later packet of that socket to whatever address it carried. A socket whose first packet went to an +// address the routing sends "direct" (a rule by address, port or sniffed protocol) carried its packets +// to all other public addresses past the VPN as well: a STUN request of the same socket showed a +// foreign server the real address of the device. Now, before a packet goes to an address other than +// the one the connection was routed by, freedom asks the router where a connection to that address +// would go (same inbound, source and sniffed protocol; no domain: it belonged to the first packet) and +// drops the packet unless the answer is this outbound. +type udpRouteGuard struct { + router routing.Router + ohm outbound.Manager + inbound *session.Inbound + content *session.Content + tag string + first net.Destination + checked map[net.Destination]bool +} + +// udpRouteGuardCache: destinations a connection remembers the answer for (a DHT socket talks to +// thousands); full - forgotten and asked again. +const udpRouteGuardCache = 256 + +func (h *Handler) newUDPRouteGuard(ctx context.Context, ob *session.Outbound) *udpRouteGuard { + if h.instance == nil || h.usesDialerProxy || h.config.DestinationOverride != nil { + // Not the final outbound, or the address is the one of the config: packets do not choose it. + return nil + } + router, _ := h.instance.GetFeature(routing.RouterType()).(routing.Router) + ohm, _ := h.instance.GetFeature(outbound.ManagerType()).(outbound.Manager) + if router == nil || ohm == nil { + return nil + } + guard := &udpRouteGuard{ + router: router, + ohm: ohm, + inbound: session.InboundFromContext(ctx), + tag: ob.Tag, + first: ob.OriginalTarget, + checked: make(map[net.Destination]bool), + } + if !guard.first.IsValid() { + guard.first = ob.Target + } + if content := session.ContentFromContext(ctx); content != nil { + guard.content = &session.Content{Protocol: content.Protocol, Attributes: content.Attributes} + } + return guard +} + +// allows: requested - the address as the packet carried it, dest - the address it is about to be sent +// to (the same, or the IP a domain was resolved to). +func (g *udpRouteGuard) allows(requested, dest net.Destination) bool { + if g == nil || requested == g.first || dest == g.first { + return true + } + if allowed, found := g.checked[dest]; found { + return allowed + } + var tag string + route, err := g.router.PickRoute(&routing_session.Context{ + Inbound: g.inbound, + Outbound: &session.Outbound{OriginalTarget: dest, Target: dest}, + Content: g.content, + }) + if err == nil { + tag = route.GetOutboundTag() + } else if handler := g.ohm.GetDefaultHandler(); handler != nil { + tag = handler.Tag() + } + allowed := tag == g.tag + if !allowed { + errors.LogInfo(context.Background(), "UDP to ", dest, " is routed to [", tag, "], not to [", g.tag, "] the connection to ", g.first, " took: dropped") + } + if len(g.checked) >= udpRouteGuardCache { + clear(g.checked) + } + g.checked[dest] = allowed + return allowed +} + func buildFinalRule(config *FinalRuleConfig) (*FinalRule, error) { rule := &FinalRule{ action: config.GetAction(), @@ -404,6 +489,9 @@ } } else { writer = NewPacketWriter(conn, h, defaultRule, UDPOverride, destination, outGateway) + if packetWriter, ok := writer.(*PacketWriter); ok { + packetWriter.RouteGuard = h.newUDPRouteGuard(ctx, ob) + } if h.config.Noises != nil { errors.LogDebug(ctx, "NOISE", h.config.Noises) writer = &NoisePacketWriter{ @@ -575,6 +663,8 @@ // So, cache and keep the resolve result ResolvedUDPAddr *utils.TypedSyncMap[string, net.Address] OutGateway net.Address + // Spanghew: see udpRouteGuard; nil - no check. + RouteGuard *udpRouteGuard } func (w *PacketWriter) WriteMultiBuffer(mb buf.MultiBuffer) error { @@ -587,6 +677,7 @@ var n int var err error if b.UDP != nil { + requested := *b.UDP if w.UDPOverride.Address != nil { b.UDP.Address = w.UDPOverride.Address } @@ -629,6 +720,10 @@ b.Release() continue } + if !w.RouteGuard.allows(requested, *b.UDP) { + b.Release() + continue + } destAddr := b.UDP.RawNetAddr() if destAddr == nil { b.Release() --- /dev/null +++ b/proxy/freedom/spanghew_udp_route_test.go @@ -0,0 +1,235 @@ +package freedom + +import ( + "context" + "testing" + "time" + + "github.com/xtls/xray-core/app/router" + "github.com/xtls/xray-core/common" + "github.com/xtls/xray-core/common/buf" + "github.com/xtls/xray-core/common/geodata" + "github.com/xtls/xray-core/common/net" + "github.com/xtls/xray-core/common/serial" + "github.com/xtls/xray-core/common/session" + "github.com/xtls/xray-core/features/outbound" + "github.com/xtls/xray-core/transport" + "github.com/xtls/xray-core/transport/internet" +) + +type spanghewHandler struct{ tag string } + +func (h *spanghewHandler) Start() error { return nil } +func (h *spanghewHandler) Close() error { return nil } +func (h *spanghewHandler) Tag() string { return h.tag } +func (h *spanghewHandler) Dispatch(context.Context, *transport.Link) {} +func (h *spanghewHandler) SenderSettings() *serial.TypedMessage { return nil } +func (h *spanghewHandler) ProxySettings() *serial.TypedMessage { return nil } + +// The first outbound ("proxy") is the default one, as in a subscription config. +type spanghewOutbounds struct { + outbound.Manager + outbound.HandlerSelector +} + +func (spanghewOutbounds) GetDefaultHandler() outbound.Handler { return &spanghewHandler{tag: "proxy"} } + +func spanghewCIDR(ip string, prefix uint32) *geodata.IPRule { + return &geodata.IPRule{Value: &geodata.IPRule_Custom{Custom: &geodata.CIDRRule{ + Cidr: &geodata.CIDR{Ip: net.ParseAddress(ip).IP(), Prefix: prefix}, + }}} +} + +// Rules of the shape a subscription template has: torrents and "home" addresses - direct, some of +// the home addresses - to the server all the same (the rule stands earlier), the rest - the default +// outbound. +func spanghewRouter(t *testing.T) *router.Router { + t.Helper() + r := new(router.Router) + common.Must(r.Init(context.TODO(), &router.Config{ + Rule: []*router.RoutingRule{ + { + TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, + Protocol: []string{"bittorrent"}, + }, + { + TargetTag: &router.RoutingRule_Tag{Tag: "proxy"}, + Ip: []*geodata.IPRule{spanghewCIDR("203.0.113.128", 25)}, + }, + { + TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, + Ip: []*geodata.IPRule{spanghewCIDR("203.0.113.0", 24)}, + }, + { + TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, + Networks: []net.Network{net.Network_UDP}, + PortList: &net.PortList{Range: []*net.PortRange{{From: 6881, To: 6889}}}, + }, + { + TargetTag: &router.RoutingRule_Tag{Tag: "other-direct"}, + Ip: []*geodata.IPRule{spanghewCIDR("192.168.0.0", 16)}, + }, + }, + }, nil, &spanghewOutbounds{}, nil)) + return r +} + +func spanghewGuard(t *testing.T, first net.Destination, protocol string) *udpRouteGuard { + t.Helper() + return &udpRouteGuard{ + router: spanghewRouter(t), + ohm: &spanghewOutbounds{}, + inbound: &session.Inbound{Tag: "tun-in", Source: net.UDPDestination(net.ParseAddress("172.19.0.1"), 50000)}, + content: &session.Content{Protocol: protocol}, + tag: "direct", + first: first, + checked: make(map[net.Destination]bool), + } +} + +func spanghewUDP(address string, port net.Port) net.Destination { + return net.UDPDestination(net.ParseAddress(address), port) +} + +// The reported leak: the first packet of a socket goes to an address routed "direct", the STUN +// request of the same socket - to a foreign server. The second must not leave through "direct". +func TestSpanghewUDPRouteGuard(t *testing.T) { + first := spanghewUDP("203.0.113.8", 9) + guard := spanghewGuard(t, first, "") + for dest, allowed := range map[net.Destination]bool{ + first: true, + spanghewUDP("203.0.113.9", 3478): true, // another "home" address: the same route + spanghewUDP("203.0.113.8", 3478): true, // another port of the first address + spanghewUDP("162.159.207.0", 3478): false, // a foreign address: the default outbound + spanghewUDP("8.8.8.8", 53): false, + spanghewUDP("203.0.113.200", 3478): false, // "home", but the earlier rule sends it to the server + spanghewUDP("192.168.1.1", 53): false, // direct too, but through another outbound + spanghewUDP("162.159.207.0", 6881): true, // the rule by port + spanghewUDP("2606:4700::1111", 443): false, + } { + for range 2 { // the second answer is the remembered one + if got := guard.allows(dest, dest); got != allowed { + t.Errorf("%v: allowed %v, want %v", dest, got, allowed) + } + } + } +} + +// The destination the connection was routed by is not asked about: its route may have been chosen +// by a sniffed domain, which the packets do not carry. +func TestSpanghewUDPRouteGuardFirstDestination(t *testing.T) { + first := spanghewUDP("162.159.207.0", 443) + guard := spanghewGuard(t, first, "quic") + if !guard.allows(first, first) { + t.Error("the first destination is dropped") + } + if !guard.allows(net.UDPDestination(net.DomainAddress("example.com"), 443), first) { + t.Error("the first destination, resolved, is dropped") + } + if guard.allows(spanghewUDP("162.159.207.0", 3478), spanghewUDP("162.159.207.0", 3478)) { + t.Error("another port of the first address went past the routing") + } + if guard.allows(spanghewUDP("162.159.207.1", 443), spanghewUDP("162.159.207.1", 443)) { + t.Error("another address went past the routing") + } +} + +// A rule by sniffed protocol holds for the whole connection: torrents routed "direct" reach all +// their peers. +func TestSpanghewUDPRouteGuardKeepsProtocol(t *testing.T) { + guard := spanghewGuard(t, spanghewUDP("162.159.207.0", 51413), "bittorrent") + for _, dest := range []net.Destination{spanghewUDP("8.8.8.8", 6000), spanghewUDP("203.0.113.200", 6000)} { + if !guard.allows(dest, dest) { + t.Errorf("%v: a torrent peer is dropped", dest) + } + } +} + +// Without a guard (not the final outbound, an address override) nothing changes. +func TestSpanghewUDPRouteGuardAbsent(t *testing.T) { + var guard *udpRouteGuard + dest := spanghewUDP("8.8.8.8", 53) + if !guard.allows(dest, dest) { + t.Error("a nil guard drops") + } + handler := &Handler{config: &Config{}} + if handler.newUDPRouteGuard(context.Background(), &session.Outbound{Target: dest}) != nil { + t.Error("a guard without a core instance") + } +} + +func TestSpanghewUDPRouteGuardForgets(t *testing.T) { + guard := spanghewGuard(t, spanghewUDP("203.0.113.8", 9), "") + for port := 1; port <= 3*udpRouteGuardCache; port++ { + dest := spanghewUDP("8.8.8.8", net.Port(port)) + if guard.allows(dest, dest) { + t.Fatalf("%v went past the routing", dest) + } + if len(guard.checked) > udpRouteGuardCache { + t.Fatalf("remembered %d destinations", len(guard.checked)) + } + } +} + +// Through the writer itself: of two packets of one connection only the one to the address routed +// here reaches the network. +func TestSpanghewPacketWriterDropsOtherRoute(t *testing.T) { + listen := func() *net.UDPConn { + conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.IP{127, 0, 0, 1}}) + common.Must(err) + t.Cleanup(func() { conn.Close() }) + return conn + } + home, foreign, local := listen(), listen(), listen() + homeDest := net.DestinationFromAddr(home.LocalAddr()) + foreignDest := net.DestinationFromAddr(foreign.LocalAddr()) + + r := new(router.Router) + common.Must(r.Init(context.TODO(), &router.Config{ + Rule: []*router.RoutingRule{{ + TargetTag: &router.RoutingRule_Tag{Tag: "direct"}, + PortList: &net.PortList{Range: []*net.PortRange{{From: uint32(homeDest.Port), To: uint32(homeDest.Port)}}}, + }}, + }, nil, &spanghewOutbounds{}, nil)) + + writer := &PacketWriter{ + PacketConnWrapper: &internet.PacketConnWrapper{PacketConn: local, Dest: home.LocalAddr()}, + Handler: &Handler{config: &Config{}}, + UDPOverride: net.UDPDestination(nil, 0), + RouteGuard: &udpRouteGuard{ + router: r, + ohm: &spanghewOutbounds{}, + tag: "direct", + first: homeDest, + checked: make(map[net.Destination]bool), + }, + } + packet := func(dest net.Destination, payload byte) *buf.Buffer { + b := buf.New() + b.WriteByte(payload) + b.UDP = &dest + return b + } + common.Must(writer.WriteMultiBuffer(buf.MultiBuffer{ + packet(homeDest, 1), packet(foreignDest, 2), packet(homeDest, 3), + })) + + read := func(conn *net.UDPConn) []byte { + var got []byte + data := make([]byte, 16) + for { + conn.SetReadDeadline(time.Now().Add(300 * time.Millisecond)) + n, _, err := conn.ReadFrom(data) + if err != nil { + return got + } + got = append(got, data[:n]...) + } + } + if got := read(home); len(got) != 2 || got[0] != 1 || got[1] != 3 { + t.Errorf("the routed address received %v", got) + } + if got := read(foreign); len(got) != 0 { + t.Errorf("the address of another route received %v", got) + } +} ============================================================================== xray-no-shadowsocks-2022.patch --- a/infra/conf/shadowsocks.go +++ b/infra/conf/shadowsocks.go @@ -3,17 +3,30 @@ import ( "strings" - "github.com/sagernet/sing-shadowsocks/shadowaead_2022" - C "github.com/sagernet/sing/common" "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/protocol" "github.com/xtls/xray-core/common/serial" "github.com/xtls/xray-core/common/task" "github.com/xtls/xray-core/proxy/shadowsocks" - "github.com/xtls/xray-core/proxy/shadowsocks_2022" "google.golang.org/protobuf/proto" ) +// Spanghew: Shadowsocks 2022 is cut out of this build. It was the only user of +// github.com/sagernet/sing and sing-shadowsocks (GPLv3) through +// proxy/shadowsocks_2022 and common/singbridge, and the client is closed source +// (Android audit 2026-10-03, #5). Those two packages stay in the tree but nothing +// imports them, so they are not linked. The methods get a clear error instead of +// "unknown cipher method". +func isShadowsocks2022(cipher string) bool { + switch cipher { + case "2022-blake3-aes-128-gcm", "2022-blake3-aes-256-gcm", "2022-blake3-chacha20-poly1305": + return true + } + return false +} + +var errShadowsocks2022 = errors.New("Shadowsocks 2022 is not included in this build") + func cipherFromString(c string) shadowsocks.CipherType { switch strings.ToLower(c) { case "aes-128-gcm", "aead_aes_128_gcm": @@ -55,8 +68,8 @@ v.Users = v.Clients } - if C.Contains(shadowaead_2022.List, v.Cipher) { - return buildShadowsocks2022(v) + if isShadowsocks2022(v.Cipher) { + return nil, errShadowsocks2022 } config := new(shadowsocks.ServerConfig) @@ -110,72 +123,6 @@ return config, nil } -func buildShadowsocks2022(v *ShadowsocksServerConfig) (proto.Message, error) { - if len(v.Users) == 0 { - config := new(shadowsocks_2022.ServerConfig) - config.Method = v.Cipher - config.Key = v.Password - config.Network = v.NetworkList.Build() - config.Email = v.Email - return config, nil - } - - if v.Cipher == "" { - return nil, errors.New("shadowsocks 2022 (multi-user): missing server method") - } - if !strings.Contains(v.Cipher, "aes") { - return nil, errors.New("shadowsocks 2022 (multi-user): only blake3-aes-*-gcm methods are supported") - } - - if v.Users[0].Address == nil { - config := new(shadowsocks_2022.MultiUserServerConfig) - config.Method = v.Cipher - config.Key = v.Password - config.Network = v.NetworkList.Build() - - config.Users = make([]*protocol.User, len(v.Users)) - processUser := func(idx int) error { - user := v.Users[idx] - if user.Cipher != "" { - return errors.New("shadowsocks 2022 (multi-user): users must have empty method") - } - account := &shadowsocks_2022.Account{ - Key: user.Password, - } - config.Users[idx] = &protocol.User{ - Email: user.Email, - Level: uint32(user.Level), - Account: serial.ToTypedMessage(account), - } - return nil - } - if err := task.ParallelForN(len(v.Users), processUser); err != nil { - return nil, err - } - return config, nil - } - - config := new(shadowsocks_2022.RelayServerConfig) - config.Method = v.Cipher - config.Key = v.Password - config.Network = v.NetworkList.Build() - for _, user := range v.Users { - if user.Cipher != "" { - return nil, errors.New("shadowsocks 2022 (relay): users must have empty method") - } - if user.Address == nil { - return nil, errors.New("shadowsocks 2022 (relay): all users must have relay address") - } - config.Destinations = append(config.Destinations, &shadowsocks_2022.RelayDestination{ - Key: user.Password, - Email: user.Email, - Address: user.Address.Build(), - Port: uint32(user.Port), - }) - } - return config, nil -} - type ShadowsocksServerTarget struct { Address *Address `json:"address"` Port uint16 `json:"port"` @@ -214,32 +161,10 @@ return nil, errors.New(`Shadowsocks settings: "servers" should have one and only one member. Multiple endpoints in "servers" should use multiple Shadowsocks outbounds and routing balancer instead`) } - if len(v.Servers) == 1 { - server := v.Servers[0] - if C.Contains(shadowaead_2022.List, server.Cipher) { - if server.Address == nil { - return nil, errors.New("Shadowsocks server address is not set.") - } - if server.Port == 0 { - return nil, errors.New("Invalid Shadowsocks port.") - } - if server.Password == "" { - return nil, errors.New("Shadowsocks password is not specified.") - } - - config := new(shadowsocks_2022.ClientConfig) - config.Address = server.Address.Build() - config.Port = uint32(server.Port) - config.Method = server.Cipher - config.Key = server.Password - return config, nil - } - } - config := new(shadowsocks.ClientConfig) for _, server := range v.Servers { - if C.Contains(shadowaead_2022.List, server.Cipher) { - return nil, errors.New("Shadowsocks 2022 accept no multi servers") + if isShadowsocks2022(server.Cipher) { + return nil, errShadowsocks2022 } if server.Address == nil { return nil, errors.New("Shadowsocks server address is not set.") --- a/main/commands/all/api/inbound_user_add.go +++ b/main/commands/all/api/inbound_user_add.go @@ -13,7 +13,6 @@ "github.com/xtls/xray-core/infra/conf" "github.com/xtls/xray-core/infra/conf/serial" "github.com/xtls/xray-core/proxy/shadowsocks" - "github.com/xtls/xray-core/proxy/shadowsocks_2022" "github.com/xtls/xray-core/proxy/trojan" vlessin "github.com/xtls/xray-core/proxy/vless/inbound" vmessin "github.com/xtls/xray-core/proxy/vmess/inbound" @@ -86,8 +85,6 @@ return ty.Users case *shadowsocks.ServerConfig: return ty.Users - case *shadowsocks_2022.MultiUserServerConfig: - return ty.Users default: fmt.Println("unsupported inbound type") } ============================================================================== xray-system-dialer-hold.patch --- a/transport/internet/dialer.go +++ b/transport/internet/dialer.go @@ -4,6 +4,7 @@ "context" "fmt" "strings" + "sync/atomic" "github.com/xtls/xray-core/common" "github.com/xtls/xray-core/common/dice" @@ -79,12 +80,27 @@ return effectiveSystemDialer.DestIpAddress() } -var ( +// Spanghew: the DNS client and the outbound manager of the system dialer belong to the whole +// process, and every core.New replaced them - two plain variables written while the connections of +// a running core were reading them. The delay measurement of libXray creates its instance next to +// the running VPN core (Android, Windows): for a moment the direct outbounds of that core asked the +// resolver of the measurement (the names went past the tunnel) and found no outbound for +// dialerProxy, and a reader could see the client of one instance with the manager of the other. +// Now the pair is one value replaced atomically, and it can be held: while it is held, +// InitSystemDialer of another instance changes nothing. +type systemDialerState struct { dnsClient dns.Client obm outbound.Manager -) + held bool +} +var systemDialer atomic.Pointer[systemDialerState] + func LookupForIP(domain string, strategy DomainStrategy, localAddr net.Address) ([]net.IP, error) { + var dnsClient dns.Client + if state := systemDialer.Load(); state != nil { + dnsClient = state.dnsClient + } if dnsClient == nil { return nil, errors.New("DNS client not initialized").AtError() } @@ -268,6 +284,10 @@ } if len(sockopt.DialerProxy) > 0 { + var obm outbound.Manager + if state := systemDialer.Load(); state != nil { + obm = state.obm + } if obm == nil { return nil, errors.New("there is no outbound manager for dialerProxy").AtError() } @@ -282,6 +302,33 @@ } func InitSystemDialer(dc dns.Client, om outbound.Manager) { - dnsClient = dc - obm = om + for { + current := systemDialer.Load() + if current != nil && current.held { + return + } + if systemDialer.CompareAndSwap(current, &systemDialerState{dnsClient: dc, obm: om}) { + return + } + } } + +// HoldSystemDialer (Spanghew): the dialer works with dc and om until ReleaseSystemDialer, whatever +// instances are created meanwhile. +func HoldSystemDialer(dc dns.Client, om outbound.Manager) { + systemDialer.Store(&systemDialerState{dnsClient: dc, obm: om, held: true}) +} + +// ReleaseSystemDialer (Spanghew): the next InitSystemDialer takes effect again; until then the +// dialer keeps what it was held with. +func ReleaseSystemDialer() { + for { + current := systemDialer.Load() + if current == nil || !current.held { + return + } + if systemDialer.CompareAndSwap(current, &systemDialerState{dnsClient: current.dnsClient, obm: current.obm}) { + return + } + } +} --- /dev/null +++ b/transport/internet/spanghew_system_dialer_test.go @@ -0,0 +1,147 @@ +package internet + +import ( + "context" + "sync" + "testing" + + "github.com/xtls/xray-core/common/net" + "github.com/xtls/xray-core/features/dns" + "github.com/xtls/xray-core/features/outbound" + "github.com/xtls/xray-core/transport" +) + +// Two DNS clients of different types, as the running core (app/dns) and a measurement instance +// (the default local client) have: a reader must never see a mix of the two. +type spanghewDNSOne struct{ ip net.IP } + +func (*spanghewDNSOne) Type() interface{} { return dns.ClientType() } +func (*spanghewDNSOne) Start() error { return nil } +func (*spanghewDNSOne) Close() error { return nil } +func (c *spanghewDNSOne) LookupIP(string, dns.IPOption) ([]net.IP, uint32, error) { + return []net.IP{c.ip}, 1, nil +} + +type spanghewDNSTwo struct { + answers map[string]net.IP +} + +func (spanghewDNSTwo) Type() interface{} { return dns.ClientType() } +func (spanghewDNSTwo) Start() error { return nil } +func (spanghewDNSTwo) Close() error { return nil } +func (c spanghewDNSTwo) LookupIP(domain string, _ dns.IPOption) ([]net.IP, uint32, error) { + return []net.IP{c.answers[domain]}, 1, nil +} + +type spanghewOutbounds struct { + outbound.Manager + handler outbound.Handler +} + +func (m *spanghewOutbounds) GetHandler(string) outbound.Handler { return m.handler } + +type spanghewOutbound struct{ outbound.Handler } + +func (spanghewOutbound) Dispatch(context.Context, *transport.Link) {} + +func spanghewSaveSystemDialer(t *testing.T) { + t.Helper() + saved := systemDialer.Load() + t.Cleanup(func() { systemDialer.Store(saved) }) +} + +func spanghewLookup(t *testing.T) string { + t.Helper() + ips, err := LookupForIP("example.invalid", DomainStrategy_USE_IP, nil) + if err != nil || len(ips) != 1 { + t.Fatalf("lookup: %v %v", ips, err) + } + return ips[0].String() +} + +// While the dialer is held (by the running core), the InitSystemDialer of another instance (a +// measurement next to it) changes neither the DNS client nor the outbound manager. +func TestSpanghewSystemDialerHeld(t *testing.T) { + spanghewSaveSystemDialer(t) + running := &spanghewDNSOne{ip: net.IP{203, 0, 113, 1}} + measuring := spanghewDNSTwo{answers: map[string]net.IP{"example.invalid": {203, 0, 113, 2}}} + proxied := &spanghewOutbounds{handler: spanghewOutbound{}} + dialThroughProxy := func() error { + conn, err := DialSystem(context.Background(), net.TCPDestination(net.DomainAddress("example.invalid"), 443), + &SocketConfig{DialerProxy: "next"}) + if err == nil { + conn.Close() + } + return err + } + + InitSystemDialer(measuring, nil) + if got := spanghewLookup(t); got != "203.0.113.2" { + t.Fatalf("not held: lookup through %s", got) + } + if dialThroughProxy() == nil { + t.Fatal("not held: dialerProxy found an outbound without a manager") + } + + HoldSystemDialer(running, proxied) + InitSystemDialer(measuring, nil) + if got := spanghewLookup(t); got != "203.0.113.1" { + t.Fatalf("held: another instance took the DNS client, lookup through %s", got) + } + if err := dialThroughProxy(); err != nil { + t.Fatalf("held: another instance took the outbound manager: %v", err) + } + + // Released: the dialer keeps what it had until the next instance is created. + ReleaseSystemDialer() + if got := spanghewLookup(t); got != "203.0.113.1" { + t.Fatalf("released: lookup through %s", got) + } + InitSystemDialer(measuring, nil) + if got := spanghewLookup(t); got != "203.0.113.2" { + t.Fatalf("after release: lookup through %s", got) + } + ReleaseSystemDialer() // not held: nothing changes + if got := spanghewLookup(t); got != "203.0.113.2" { + t.Fatalf("release of a dialer that is not held: lookup through %s", got) + } +} + +// Instances come and go while connections resolve names: every lookup gets the answer of one whole +// client (run with -race to see the readers and the writer). +func TestSpanghewSystemDialerSwappedUnderLookups(t *testing.T) { + spanghewSaveSystemDialer(t) + one := &spanghewDNSOne{ip: net.IP{203, 0, 113, 1}} + two := spanghewDNSTwo{answers: map[string]net.IP{"example.invalid": {203, 0, 113, 2}}} + InitSystemDialer(one, nil) + + stop := make(chan struct{}) + var readers sync.WaitGroup + for range 4 { + readers.Add(1) + go func() { + defer readers.Done() + for { + select { + case <-stop: + return + default: + } + ips, err := LookupForIP("example.invalid", DomainStrategy_USE_IP, nil) + if err != nil || len(ips) != 1 || (ips[0].String() != "203.0.113.1" && ips[0].String() != "203.0.113.2") { + t.Errorf("lookup: %v %v", ips, err) + return + } + } + }() + } + for i := range 20000 { + if i%2 == 0 { + InitSystemDialer(two, nil) + } else { + InitSystemDialer(one, nil) + } + } + close(stop) + readers.Wait() +} ============================================================================== xray-tun-handshake-log.patch --- a/proxy/tun/stack_gvisor.go +++ b/proxy/tun/stack_gvisor.go @@ -74,7 +74,11 @@ // Perform a TCP three-way handshake. ep, err := r.CreateEndpoint(&wq) if err != nil { - errors.LogError(t.ctx, err.String()) + // Spanghew: the app did not finish the handshake of its own connection (it reset or + // dropped it) - an everyday event, not a fault of the core. Upstream logged it as an + // error: dozens of "connection reset by peer" lines a day in the log the user sees, + // pushing the lines that matter out of it. + errors.LogInfo(t.ctx, "connection from the tunnel not established: ", err.String()) r.Complete(true) return } ============================================================================== xray-tun-owner-check.patch --- a/proxy/tun/handler.go +++ b/proxy/tun/handler.go @@ -3,6 +3,7 @@ import ( "context" "net/netip" + "runtime" "strings" "syscall" @@ -171,6 +172,10 @@ return } source := net.DestinationFromAddr(remote) + if !ownerAllowed(source, destination) { + errors.LogInfo(t.ctx, "blocked connection of an excluded or unknown app from ", source, " to ", destination) + return + } if t.uplinkCounter != nil || t.downlinkCounter != nil { conn = &stat.CounterConnection{ Connection: conn, @@ -239,3 +244,27 @@ return t, err })) } + +// ownerAllowed (Spanghew): пускать ли соединение в туннель, по владельцу — uid приложения +// от net.FindProcess (на Android его ищет приложение через libXray.RegisterProcessFinder и +// ConnectivityManager.getConnectionOwnerUid). Отрицательный uid — «не пускать»: владелец не найден +// (так бывает у сокета, привязанного к tun0 через SO_BINDTODEVICE, — приложение из исключений VPN +// так узнало бы адрес сервера) или приложение исключено из VPN. Поиск не зарегистрирован или не +// смог спросить — пропускаем, как без заплатки. Адреса — исходные, до sniffing. +// Только Android: на macOS и Windows FindProcess отдаёт PID ≥ 0 или ошибку (соединение всё равно пропускается), +// но на каждое соединение TUN перебирает сокеты всех процессов или таблицы соединений Windows (аудит M24d, w7). +func ownerAllowed(source, destination net.Destination) bool { + if runtime.GOOS != "android" { + return true + } + if !source.Address.Family().IsIP() || !destination.Address.Family().IsIP() { + return true + } + network := "tcp" + if destination.Network == net.Network_UDP { + network = "udp" + } + uid, _, _, err := net.FindProcess(network, source.Address.IP().String(), uint16(source.Port), + destination.Address.IP().String(), uint16(destination.Port)) + return err != nil || uid >= 0 +} ============================================================================== xray-tun-udp-scope.patch --- a/proxy/tun/udp_fullcone.go +++ b/proxy/tun/udp_fullcone.go @@ -3,6 +3,7 @@ import ( "context" "io" + "net/netip" "sync" "time" @@ -20,7 +21,7 @@ type udpConnectionHandler struct { sync.RWMutex - udpConns map[net.Destination]*udpConn + udpConns map[udpConnKey]*udpConn handleConnection func(conn net.Conn, dest net.Destination) writePacket func(data []byte, src net.Destination, dst net.Destination) error @@ -28,7 +29,7 @@ func newUdpConnectionHandler(handleConnection func(conn net.Conn, dest net.Destination), writePacket func(data []byte, src net.Destination, dst net.Destination) error) *udpConnectionHandler { handler := &udpConnectionHandler{ - udpConns: make(map[net.Destination]*udpConn), + udpConns: make(map[udpConnKey]*udpConn), handleConnection: handleConnection, writePacket: writePacket, } @@ -36,11 +37,78 @@ return handler } +// Spanghew: the route of a connection is chosen by its first packet. Upstream keyed connections by +// the source alone, so a socket whose first packet went to a local address (routed "direct" or to +// the local-network block) carried all its later packets to public addresses the same way - past +// the VPN. Local and public destinations of one socket are separate connections now; FullCone NAT +// stays within each. +type udpConnKey struct { + src net.Destination + local bool +} + +func udpKey(src, dst net.Destination) udpConnKey { + return udpConnKey{src: src, local: isLocalDestination(dst)} +} + +// localNetworks: geoip:private of the app's geobase (v2fly/geoip) - what subscription templates +// route "direct". Private (RFC 1918, fc00::/7), loopback, link-local, multicast and broadcast +// addresses are in it, and so are the reserved ranges a public host never answers from +// (0.0.0.0/8, 100.64.0.0/10, 192.0.2.0/24, 198.18.0.0/15, 240.0.0.0/4 ...): a first packet there +// would take the "direct" route as well. The same list - udp_pretend_scope of the hev lwIP patch; +// the app's test compares both with the geobase. +var localNetworks = func() []netip.Prefix { + var networks []netip.Prefix + for _, network := range []string{ + "0.0.0.0/8", + "10.0.0.0/8", + "100.64.0.0/10", + "127.0.0.0/8", + "169.254.0.0/16", + "172.16.0.0/12", + "192.0.0.0/24", + "192.0.2.0/24", + "192.88.99.0/24", + "192.168.0.0/16", + "198.18.0.0/15", + "198.51.100.0/24", + "203.0.113.0/24", + "224.0.0.0/3", + "::/127", + "fc00::/7", + "fe80::/10", + "ff00::/8", + } { + networks = append(networks, netip.MustParsePrefix(network)) + } + return networks +}() + +// isLocalDestination: an address of localNetworks (IPv4 inside IPv6 "::ffff:..." - as IPv4, the way +// the router reads it). +func isLocalDestination(dst net.Destination) bool { + if dst.Address == nil || !dst.Address.Family().IsIP() { + return false + } + ip, ok := netip.AddrFromSlice(dst.Address.IP()) + if !ok { + return false + } + ip = ip.Unmap() + for _, network := range localNetworks { + if network.Contains(ip) { + return true + } + } + return false +} + // HandlePacket handles UDP packets coming from tun, to forward to the dispatcher // this custom handler support FullCone NAT of returning packets, binding connection only by the source addr:port func (u *udpConnectionHandler) HandlePacket(src net.Destination, dst net.Destination, data []byte) { + key := udpKey(src, dst) u.RLock() - conn, found := u.udpConns[src] + conn, found := u.udpConns[key] if found { select { case conn.egress <- &packet{ @@ -58,11 +126,11 @@ u.Lock() defer u.Unlock() - conn, found = u.udpConns[src] + conn, found = u.udpConns[key] if !found { egress := make(chan *packet, 1024) conn = &udpConn{handler: u, egress: egress, src: src, dst: dst} - u.udpConns[src] = conn + u.udpConns[key] = conn go u.handleConnection(conn, dst) } @@ -78,11 +146,11 @@ } } -func (u *udpConnectionHandler) connectionFinished(src net.Destination) { +func (u *udpConnectionHandler) connectionFinished(key udpConnKey) { u.Lock() - conn, found := u.udpConns[src] + conn, found := u.udpConns[key] if found { - delete(u.udpConns, src) + delete(u.udpConns, key) close(conn.egress) } u.Unlock() @@ -161,7 +229,7 @@ } func (c *udpConn) Close() error { - c.handler.connectionFinished(c.src) + c.handler.connectionFinished(udpKey(c.src, c.dst)) return nil } --- /dev/null +++ b/proxy/tun/spanghew_udp_scope_test.go @@ -0,0 +1,143 @@ +package tun + +import ( + "sync" + "testing" + + "github.com/xtls/xray-core/common/net" +) + +// Spanghew: one socket, first packet to the local network, the next to a public address - two +// connections, each routed by its own first packet; packets to other public addresses stay in the +// public one (FullCone). +func TestSpanghewUDPLocalAndPublicAreSeparateConnections(t *testing.T) { + var ( + mu sync.Mutex + opened []net.Destination + conns []net.Conn + wait sync.WaitGroup + ) + handler := newUdpConnectionHandler(func(conn net.Conn, dest net.Destination) { + mu.Lock() + opened = append(opened, dest) + conns = append(conns, conn) + mu.Unlock() + wait.Done() + }, func([]byte, net.Destination, net.Destination) error { return nil }) + + src := net.UDPDestination(net.ParseAddress("172.19.0.1"), 50000) + router := net.UDPDestination(net.ParseAddress("192.168.1.1"), 3478) + public := net.UDPDestination(net.ParseAddress("1.1.1.1"), 3478) + other := net.UDPDestination(net.ParseAddress("8.8.8.8"), 3478) + + wait.Add(2) + handler.HandlePacket(src, router, []byte{1}) + handler.HandlePacket(src, public, []byte{2}) + handler.HandlePacket(src, other, []byte{3}) + handler.HandlePacket(src, router, []byte{4}) + wait.Wait() + + mu.Lock() + defer mu.Unlock() + if len(opened) != 2 { + t.Fatalf("connections: %v", opened) + } + seen := map[net.Destination]bool{opened[0]: true, opened[1]: true} + if !seen[router] || !seen[public] { + t.Fatalf("connections routed by: %v", opened) + } + if len(handler.udpConns) != 2 { + t.Fatalf("connections kept: %d", len(handler.udpConns)) + } + for _, conn := range conns { + _ = conn.Close() + } + if len(handler.udpConns) != 0 { + t.Fatalf("connections left after close: %d", len(handler.udpConns)) + } +} + +func TestSpanghewLocalDestinations(t *testing.T) { + for address, local := range map[string]bool{ + "10.1.2.3": true, + "172.16.0.9": true, + "192.168.1.255": true, + "169.254.10.10": true, + "127.0.0.1": true, + "224.0.0.251": true, + "255.255.255.255": true, + "0.1.2.3": true, + "100.64.0.1": true, + "100.127.255.255": true, + "192.0.0.9": true, + "192.0.2.1": true, + "192.88.99.1": true, + "198.18.0.1": true, + "198.19.255.255": true, + "198.51.100.7": true, + "203.0.113.7": true, + "240.0.0.1": true, + "::": true, + "::1": true, + "fe80::1": true, + "fd00::1": true, + "ff02::fb": true, + "::ffff:192.0.2.1": true, + "::ffff:192.168.1.1": true, + "1.1.1.1": false, + "100.63.255.255": false, + "100.128.0.1": false, + "192.0.1.1": false, + "198.17.255.255": false, + "198.20.0.1": false, + "223.255.255.255": false, + "::2": false, + "::ffff:1.1.1.1": false, + "2606:4700::1111": false, + "2001:db8::1": false, + } { + if got := isLocalDestination(net.UDPDestination(net.ParseAddress(address), 53)); got != local { + t.Errorf("%s: local %v, want %v", address, got, local) + } + } + if isLocalDestination(net.UDPDestination(net.ParseAddress("example.com"), 53)) { + t.Error("a domain is not a local destination") + } +} + +// The reported leak: the first packet of a socket goes to a reserved address that geoip:private holds +// and RFC 1918 does not; the socket's packets to public addresses are another connection. +func TestSpanghewUDPReservedAndPublicAreSeparateConnections(t *testing.T) { + var ( + mu sync.Mutex + opened []net.Destination + conns []net.Conn + wait sync.WaitGroup + ) + handler := newUdpConnectionHandler(func(conn net.Conn, dest net.Destination) { + mu.Lock() + opened = append(opened, dest) + conns = append(conns, conn) + mu.Unlock() + wait.Done() + }, func([]byte, net.Destination, net.Destination) error { return nil }) + + src := net.UDPDestination(net.ParseAddress("172.19.0.1"), 50000) + reserved := net.UDPDestination(net.ParseAddress("192.0.2.1"), 9) + public := net.UDPDestination(net.ParseAddress("1.1.1.1"), 3478) + + wait.Add(2) + handler.HandlePacket(src, reserved, []byte{1}) + handler.HandlePacket(src, public, []byte{2}) + wait.Wait() + + mu.Lock() + defer mu.Unlock() + seen := map[net.Destination]bool{opened[0]: true, opened[1]: true} + if len(opened) != 2 || !seen[reserved] || !seen[public] { + t.Fatalf("connections routed by: %v", opened) + } + for _, conn := range conns { + _ = conn.Close() + } +} ============================================================================== xray-tun-windows-default-route.patch --- a/proxy/tun/handler.go +++ b/proxy/tun/handler.go @@ -5,6 +5,7 @@ import ( "net/netip" "runtime" "strings" + "sync" "syscall" "github.com/xtls/xray-core/common" @@ -81,6 +82,31 @@ func (t *Handler) Init(ctx context.Context, pm policy. } return nil +} + +// Spanghew: see Start. +var interfaceControllerOnce sync.Once + +func bindToOutboundInterface(network, address string, c syscall.RawConn) error { + addrPort, _ := netip.ParseAddrPort(address) + // skip loopback + if addrPort.Addr().IsLoopback() || strings.HasPrefix(strings.ToLower(address), "localhost:") { + return nil + } + iface := updater.Get() + if iface == nil { + // Spanghew: no interface to bind to (no network besides the TUN). Upstream left + // the socket unbound: it followed the default route back into the TUN, and the + // core talked to itself in a loop. No network - the dial fails: the dialer stops + // on this error only (transport/internet/system_dialer.go), any other it logs. + return internet.ErrControllerRefused + } + return c.Control(func(fd uintptr) { + err := setinterface(network, address, fd, iface) + if err != nil { + errors.LogInfoInner(context.Background(), err, "[tun] falied to set interface") + } + }) } func (t *Handler) Start() error { @@ -101,23 +127,11 @@ func (t *Handler) Start() error { } updater = &InterfaceUpdater{tunIndex: tunIndex, fixedName: t.config.AutoOutboundsInterface} updater.Update() - internet.RegisterDialerController(func(network, address string, c syscall.RawConn) error { - iface := updater.Get() - if iface == nil { - errors.LogInfo(context.Background(), "[tun] falied to set interface > iface == nil") - return nil - } - return c.Control(func(fd uintptr) { - addrPort, _ := netip.ParseAddrPort(address) - // skip loopback - if addrPort.Addr().IsLoopback() || strings.HasPrefix(strings.ToLower(address), "localhost:") { - return - } - err := setinterface(network, address, fd, iface) - if err != nil { - errors.LogInfoInner(context.Background(), err, "[tun] falied to set interface") - } - }) + // Spanghew: one controller per process. It reads the package-level updater, and + // upstream appended another copy on every start of the inbound - a long-lived + // process (the Windows service) bound each socket once per past connection. + interfaceControllerOnce.Do(func() { + internet.RegisterDialerController(bindToOutboundInterface) }) } --- a/proxy/tun/tun_windows.go +++ b/proxy/tun/tun_windows.go @@ -205,6 +205,12 @@ if err != nil { return err } + // Spanghew: the first Update() ran before the TUN had its routes (handler.go), and what + // changed between it and the two registrations above nobody reported: a default route + // that appeared in that gap (the network coming up while the tunnel starts) left the + // updater without an interface, and the controller refused every dial until the next + // change of the network. Look once more, now that changes are reported. + updater.Update() } return nil } @@ -222,6 +228,14 @@ } if t.cbi != nil { t.cbi.Unregister() + } + // Spanghew: туннеля больше нет — и привязки к его прежней сетевой карте тоже. Контроллер сокетов один на + // процесс (handler.go) и читает updater, а обновлять его после Close некому: служба Windows живёт дальше, + // и замеры задержки шли бы через карту прошлого туннеля (сменили Wi-Fi на кабель — «нет ответа») или не + // шли вовсе (туннель остановили без сети). Номер 0 для IP_UNICAST_IF — «без привязки», маршрут выбирает + // Windows. Следующий Start ставит свой updater. + if updater != nil { + updater = &InterfaceUpdater{iface: &net.Interface{}} } if t.luid != 0 { t.luid.FlushRoutes(windows.AF_INET) @@ -349,45 +363,47 @@ return net.InterfaceByName(fixedName) } - r, err := winipcfg.GetIPForwardTable2(windows.AF_UNSPEC) - if err != nil { - return nil, err + // Spanghew: интерфейс, через который система и так ходит в интернет, — с маршрутом по умолчанию + // и наименьшей суммой метрик, кроме самого TUN. Ядро с 26.9.8 тоже выбирает по маршрутам, но Wi-Fi + // берёт раньше кабеля с меньшей метрикой и не различает IPv4 и IPv6. Нет такого интерфейса — nil: + // контроллер сокетов (handler.go) откажет в соединении, а не пустит его обратно в TUN. + index := defaultRouteInterfaceIndex(tunIndex) + if index == 0 { + return nil, nil } - lowestMetric := ^uint32(0) - index := uint32(0) - lowestMetricWifi := ^uint32(0) - indexWifi := uint32(0) - for i := range r { - if r[i].DestinationPrefix.PrefixLength != 0 || r[i].InterfaceIndex == uint32(tunIndex) { - continue + return net.InterfaceByIndex(index) +} + +// defaultRouteInterfaceIndex (Spanghew): номер подключённого интерфейса с маршрутом по умолчанию +// и наименьшей суммой метрик маршрута и интерфейса (как выбирает сама Windows); 0 — такого нет. +// Сначала IPv4 (0.0.0.0/0); нет — IPv6 (::/0): сеть только с IPv6. +func defaultRouteInterfaceIndex(tunIndex int) int { + for _, family := range []winipcfg.AddressFamily{windows.AF_INET, windows.AF_INET6} { + if index := defaultRouteInterfaceIndexOf(family, tunIndex); index != 0 { + return index } - ifrow, err := r[i].InterfaceLUID.Interface() - if err != nil || ifrow.OperStatus != winipcfg.IfOperStatusUp { + } + return 0 +} + +func defaultRouteInterfaceIndexOf(family winipcfg.AddressFamily, tunIndex int) int { + routes, err := winipcfg.GetIPForwardTable2(family) + if err != nil { + return 0 + } + best, bestMetric := 0, uint64(0) + for _, route := range routes { + if route.DestinationPrefix.PrefixLength != 0 || int(route.InterfaceIndex) == tunIndex { continue } - - iface, err := r[i].InterfaceLUID.IPInterface(windows.AF_INET) - if err != nil { - iface, err = r[i].InterfaceLUID.IPInterface(windows.AF_INET6) - if err != nil { - continue - } - } - - if ifrow.Type == windows.IF_TYPE_IEEE80211 { - if r[i].Metric+iface.Metric < lowestMetricWifi { - lowestMetricWifi = r[i].Metric + iface.Metric - indexWifi = r[i].InterfaceIndex - } + ipif, err := route.InterfaceLUID.IPInterface(family) + if err != nil || !ipif.Connected { continue } - if r[i].Metric+iface.Metric < lowestMetric { - lowestMetric = r[i].Metric + iface.Metric - index = r[i].InterfaceIndex + metric := uint64(route.Metric) + uint64(ipif.Metric) + if best == 0 || metric < bestMetric { + best, bestMetric = int(route.InterfaceIndex), metric } } - if indexWifi != 0 { - index = indexWifi - } - return net.InterfaceByIndex(int(index)) + return best } --- a/transport/internet/system_dialer.go +++ b/transport/internet/system_dialer.go @@ -2,6 +2,7 @@ import ( import ( "context" + goerrors "errors" "sync" "syscall" "time" @@ -18,6 +19,13 @@ var ( effectiveSystemDialer SystemDialer = &DefaultSystemDialer{} ) +// ErrControllerRefused (Spanghew): a dialer controller returns it to fail the dial. Upstream only +// logs what a controller returns and dials on, so a controller could not stop a socket it had not +// prepared: on Windows a socket left unbound with no network besides the TUN followed the default +// route back into the TUN (the controller is in proxy/tun/handler.go). Any other error of a +// controller is logged and ignored, as before. +var ErrControllerRefused = goerrors.New("dial refused by a dialer controller") + type SystemDialer interface { Dial(ctx context.Context, source net.Address, destination net.Destination, sockopt *SocketConfig) (net.Conn, error) DestIpAddress() net.IP @@ -71,6 +79,9 @@ func (d *DefaultSystemDialer) Dial(ctx context.Context lc.Control = func(network, address string, c syscall.RawConn) error { for _, ctl := range Controllers { if err := ctl(network, address, c); err != nil { + if goerrors.Is(err, ErrControllerRefused) { + return err + } errors.LogInfoInner(ctx, err, "failed to apply external controller") } } @@ -128,6 +139,9 @@ func (d *DefaultSystemDialer) Dial(ctx context.Context dialer.Control = func(network, address string, c syscall.RawConn) error { for _, ctl := range Controllers { if err := ctl(network, address, c); err != nil { + if goerrors.Is(err, ErrControllerRefused) { + return err + } errors.LogInfoInner(ctx, err, "failed to apply external controller") } } --- /dev/null +++ b/transport/internet/spanghew_controller_refused_test.go @@ -0,0 +1,72 @@ +package internet + +import ( + "context" + goerrors "errors" + "syscall" + "testing" + + "github.com/xtls/xray-core/common/net" +) + +// Spanghew: a controller that returns ErrControllerRefused fails the dial - TCP and UDP; any other +// error of a controller is logged and the dial goes on, as upstream. +func TestSpanghewControllerRefusesDial(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer listener.Close() + go func() { + for { + conn, err := listener.Accept() + if err != nil { + return + } + conn.Close() + } + }() + port := net.Port(listener.Addr().(*net.TCPAddr).Port) + destinations := []net.Destination{ + net.TCPDestination(net.LocalHostIP, port), + net.UDPDestination(net.LocalHostIP, port), + } + + ControllersLock.Lock() + saved := Controllers + ControllersLock.Unlock() + defer func() { + ControllersLock.Lock() + Controllers = saved + ControllersLock.Unlock() + }() + set := func(result error) { + ControllersLock.Lock() + Controllers = []func(network, address string, c syscall.RawConn) error{ + func(string, string, syscall.RawConn) error { return result }, + } + ControllersLock.Unlock() + } + dialer := &DefaultSystemDialer{} + + set(ErrControllerRefused) + for _, dest := range destinations { + conn, err := dialer.Dial(context.Background(), nil, dest, nil) + if conn != nil { + conn.Close() + } + if !goerrors.Is(err, ErrControllerRefused) { + t.Errorf("%s: refused controller, dial error: %v", dest, err) + } + } + + set(goerrors.New("could not prepare the socket")) + for _, dest := range destinations { + conn, err := dialer.Dial(context.Background(), nil, dest, nil) + if err != nil { + t.Errorf("%s: another controller error must not stop the dial: %v", dest, err) + continue + } + conn.Close() + } +} ============================================================================== xray-windows-long-process-path.patch --- a/common/net/find_process_windows.go +++ b/common/net/find_process_windows.go @@ -113,6 +113,11 @@ return 0, "", "", err } NameWithPath, err := getExecPathFromPID(pid) + // Spanghew: у программы, запущенной по короткому пути 8.3 (C:\PROGRA~1\...), Windows отдаёт путь с короткими + // именами, а правила process — по длинным: без этого выбранная программа шла бы мимо VPN. + if err == nil && strings.Contains(NameWithPath, "~") { + NameWithPath = longPath(NameWithPath) + } NameWithPath = filepath.ToSlash(NameWithPath) // drop .exe and path @@ -222,6 +227,20 @@ return *(*uint32)(unsafe.Pointer(&b[0])) } +// Spanghew: длинная форма пути (GetLongPathNameW); не вышло — путь как есть. +func longPath(path string) string { + p, err := windows.UTF16PtrFromString(path) + if err != nil { + return path + } + buf := make([]uint16, syscall.MAX_LONG_PATH) + n, err := windows.GetLongPathName(p, &buf[0], uint32(len(buf))) + if err != nil || n == 0 || n >= uint32(len(buf)) { + return path + } + return windows.UTF16ToString(buf[:n]) +} + func getExecPathFromPID(pid uint32) (string, error) { // kernel process starts with a colon in order to distinguish with normal processes switch pid { ============================================================================== apple/xray-tun-darwin-dup-fd.patch --- a/proxy/tun/tun_darwin.go +++ b/proxy/tun/tun_darwin.go @@ -164,12 +164,21 @@ return nil, err } + // Spanghew: work on our own copy of the NetworkExtension fd. os.File closes its fd from a GC finalizer; + // after a server change (core stop + start on the same utun) the old core's file would close the fd + // number the new core reads. The copy is ours alone: Close() closes it and wakes the old reader. + own, err := unix.Dup(fd) + if err != nil { + return nil, err + } + unix.CloseOnExec(own) + return &DarwinTun{ - tunFile: os.NewFile(uintptr(fd), "utun"), + tunFile: os.NewFile(uintptr(own), "utun"), options: options, - tunFd: fd, + tunFd: own, ownsFd: false, - waitKq: newWaitKqueue(fd), + waitKq: newWaitKqueue(own), }, nil } @@ -232,11 +241,9 @@ t.waitKq.close() } routeErr := t.unsetSystemRoutes() - if t.ownsFd { - return xerrors.Combine(routeErr, t.tunFile.Close()) - } - // iOS: don't close the fd, it's owned by NetworkExtension - return routeErr + // Spanghew: the fd from NetworkExtension is a copy made in NewTun — close it too; the original stays with + // the system. + return xerrors.Combine(routeErr, t.tunFile.Close()) } func (t *DarwinTun) monitorRouteChanges() { ============================================================================== apple/xray-tun-ios-buffers.patch --- a/proxy/tun/stack_gvisor.go +++ b/proxy/tun/stack_gvisor.go @@ -25,11 +25,9 @@ tcpRXBufMinSize = tcp.MinBufferSize tcpRXBufDefSize = tcp.DefaultSendBufferSize - tcpRXBufMaxSize = 8 << 20 // 8MiB tcpTXBufMinSize = tcp.MinBufferSize tcpTXBufDefSize = tcp.DefaultReceiveBufferSize - tcpTXBufMaxSize = 6 << 20 // 6MiB ) // stackGVisor is ip stack implemented by gVisor package --- /dev/null +++ b/proxy/tun/spanghew_bufsize.go @@ -0,0 +1,20 @@ +package tun + +import "runtime" + +// Spanghew: upper limits of the per-connection TCP buffers of the gVisor +// stack. Upstream: 8 MiB receive, 6 MiB send, window auto-tuning on. The iOS +// Packet Tunnel extension lives under a ~50 MiB memory limit (iPhone 7, +// iOS 15), and a few uploads through a slow server could grow the windows past +// it - iOS then kills the extension. The TUN side is local (no round trip +// time), so smaller windows do not limit the speed. Not below the defaults +// (1 MiB both): gVisor refuses a range whose maximum is under its default, +// and the stack would not start (1.0.25 on the iPhone, 3.10). +var tcpRXBufMaxSize, tcpTXBufMaxSize = tcpBufMaxSizes(runtime.GOOS) + +func tcpBufMaxSizes(goos string) (rx, tx int) { + if goos == "ios" { + return 1 << 20, 1 << 20 + } + return 8 << 20, 6 << 20 +} --- /dev/null +++ b/proxy/tun/spanghew_bufsize_test.go @@ -0,0 +1,36 @@ +package tun + +import ( + "testing" + + "gvisor.dev/gvisor/pkg/tcpip" + "gvisor.dev/gvisor/pkg/tcpip/network/ipv4" + "gvisor.dev/gvisor/pkg/tcpip/stack" + "gvisor.dev/gvisor/pkg/tcpip/transport/tcp" +) + +// The ranges must be accepted by the gVisor stack itself, the way Start sets them. +func TestSpanghewTCPBufMaxSizes(t *testing.T) { + for _, goos := range []string{"ios", "darwin", "android", "windows", "linux"} { + rx, tx := tcpBufMaxSizes(goos) + s := stack.New(stack.Options{ + NetworkProtocols: []stack.NetworkProtocolFactory{ipv4.NewProtocol}, + TransportProtocols: []stack.TransportProtocolFactory{tcp.NewProtocol}, + }) + rxOpt := tcpip.TCPReceiveBufferSizeRangeOption{Min: tcpRXBufMinSize, Default: tcpRXBufDefSize, Max: rx} + if err := s.SetTransportProtocolOption(tcp.ProtocolNumber, &rxOpt); err != nil { + t.Errorf("%s: receive range refused: %v", goos, err) + } + txOpt := tcpip.TCPSendBufferSizeRangeOption{Min: tcpTXBufMinSize, Default: tcpTXBufDefSize, Max: tx} + if err := s.SetTransportProtocolOption(tcp.ProtocolNumber, &txOpt); err != nil { + t.Errorf("%s: send range refused: %v", goos, err) + } + s.Close() + } + if rx, tx := tcpBufMaxSizes("ios"); rx != 1<<20 || tx != 1<<20 { + t.Fatalf("ios: %d, %d", rx, tx) + } + if rx, tx := tcpBufMaxSizes("darwin"); rx != 8<<20 || tx != 6<<20 { + t.Fatalf("darwin: %d, %d", rx, tx) + } +} ============================================================================== apple/xray-tun-log-read-stop.patch --- a/proxy/tun/stack_gvisor_endpoint.go 2026-09-30 16:49:28 +++ b/proxy/tun/stack_gvisor_endpoint.go 2026-09-30 16:49:28 @@ -4,6 +4,7 @@ "context" "errors" + xerrors "github.com/xtls/xray-core/common/errors" "gvisor.dev/gvisor/pkg/tcpip" "gvisor.dev/gvisor/pkg/tcpip/header" "gvisor.dev/gvisor/pkg/tcpip/stack" @@ -126,6 +127,10 @@ } // stop dispatcher loop on any other interface failure if err != nil { + // Spanghew: the loop used to stop silently — traffic into the TUN then goes nowhere. + if ctx.Err() == nil { + xerrors.LogWarningInner(ctx, err, "[tun] packet read loop stopped") + } e.Attach(nil) return }