Changing any stats knob on the persistent backend deleted months of history. The replacement store was opened before the outgoing one was closed, so it hit the first one's flock, timed out — and the open path treated ANY error as corruption and unlinked the file. Unlink of an open file succeeds on Linux, so the new ring opened an empty database while the panel was still told the backend had not changed. The store now hands its resources over before asking for them again, and deletion is gated on an allow-list of real corruption signals; a busy, unreadable or read-only file degrades to the in-RAM ring and is left alone. The holder's reads were unguarded in a subtler way, caught only after the gate failed twice: the accessor took the read lock, returned the pointer and released it, so the call ran outside. A reader could hold a store the swap then closed and be served its empty answer — an empty page presented as data. The accessor is gone entirely, along with the possibility of handing out an unguarded reference. Readers still do not block each other; the swap now waits out reads already in flight, which is a page at most. The urltest group published its chosen node through two plain fields written by the prober and read on every dial and every panel poll — while the selector next door does the same job atomically. They are one value now, so TCP and UDP can no longer be read as a mismatched pair. Nothing had ever dialled through a group while it was probing, which is why the detector had never seen it; a test now does, and reproduces it deterministically against the old shape. Close on a group whose ticker had already stopped returned before closing its channel, and Touch would then arm a fresh loop nothing could stop. Reached by pressing Test in the panel and applying a config within the next two minutes: the orphan kept failing probes against a cancelled context and writing forged dead verdicts into the board the live generation selects from. Close is now final. The log sink held one mutex across a blocking write. Under procd stderr is a pipe, so a reader that stopped draining wedged everything that logs — engine, panel handlers, signal loop — while the process still answered a signal. It is split: a front that assembles lines and a writer that owns the destinations, joined by a bounded queue that drops and counts rather than blocking. Proven by restoring the old shape: the package deadlocks for the full ten-minute timeout, parked exactly where the field symptom said. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1138 lines
44 KiB
Go
1138 lines
44 KiB
Go
package group
|
|
|
|
import (
|
|
"context"
|
|
"net"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/sagernet/sing-box/adapter"
|
|
"github.com/sagernet/sing-box/adapter/outbound"
|
|
"github.com/sagernet/sing-box/common/interrupt"
|
|
"github.com/sagernet/sing-box/common/urltest"
|
|
C "github.com/sagernet/sing-box/constant"
|
|
"github.com/sagernet/sing-box/log"
|
|
"github.com/sagernet/sing-box/option"
|
|
"github.com/sagernet/sing/common"
|
|
"github.com/sagernet/sing/common/batch"
|
|
E "github.com/sagernet/sing/common/exceptions"
|
|
M "github.com/sagernet/sing/common/metadata"
|
|
N "github.com/sagernet/sing/common/network"
|
|
"github.com/sagernet/sing/common/x/list"
|
|
"github.com/sagernet/sing/service"
|
|
"github.com/sagernet/sing/service/pause"
|
|
)
|
|
|
|
func RegisterURLTest(registry *outbound.Registry) {
|
|
outbound.Register[option.URLTestOutboundOptions](registry, C.TypeURLTest, NewURLTest)
|
|
}
|
|
|
|
var _ adapter.OutboundGroup = (*URLTest)(nil)
|
|
|
|
type URLTest struct {
|
|
outbound.Adapter
|
|
ctx context.Context
|
|
outbound adapter.OutboundManager
|
|
connection adapter.ConnectionManager
|
|
logger log.ContextLogger
|
|
tags []string
|
|
link string
|
|
interval time.Duration
|
|
tolerance uint16
|
|
idleTimeout time.Duration
|
|
group *URLTestGroup
|
|
interruptExternalConnections bool
|
|
balancer *balancer // lx: SPEC 019 — nil for least_test (default)
|
|
// lx: health board §5.C — true when options.SelfCheck == false: the group's
|
|
// OWN probing schedule (PostStart warm-up + Touch ticker) is stood down and
|
|
// the observatory is the only thing that measures its members. Stored
|
|
// INVERTED so the zero value keeps today's behaviour for every construction
|
|
// path that does not go through NewURLTest (hand-built groups in tests).
|
|
// See option.URLTestOutboundOptions.SelfCheck for the full reasoning.
|
|
selfCheckDisabled bool
|
|
// lx: health board §5.C — the RUNTIME half of the same question, read from
|
|
// the context registry at construction. selfCheckDisabled above says "this
|
|
// group is not used by the config"; this says "this group cannot be reached
|
|
// right now", which changes while the box runs and is therefore asked
|
|
// afresh at every scheduled check rather than stored. nil = no gate.
|
|
// See urltest.ProbeGate for why the two must stay separate.
|
|
probeGate urltest.ProbeGate
|
|
}
|
|
|
|
func NewURLTest(ctx context.Context, router adapter.Router, logger log.ContextLogger, tag string, options option.URLTestOutboundOptions) (adapter.Outbound, error) {
|
|
// lx: SPEC 019 v2 — round_robin balancer; nil balancer keeps legacy least_test behaviour.
|
|
balancer, err := newBalancer(options)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if balancer != nil && options.Tolerance != 0 {
|
|
logger.Warn("urltest: tolerance is ignored in round_robin mode; use balancer.pool_tolerance")
|
|
}
|
|
if balancer == nil && options.Balancer != nil {
|
|
return nil, E.New("urltest: balancer is only valid with mode: round_robin")
|
|
}
|
|
outbound := &URLTest{
|
|
Adapter: outbound.NewAdapter(C.TypeURLTest, tag, []string{N.NetworkTCP, N.NetworkUDP}, options.Outbounds),
|
|
ctx: ctx,
|
|
outbound: service.FromContext[adapter.OutboundManager](ctx),
|
|
connection: service.FromContext[adapter.ConnectionManager](ctx),
|
|
logger: logger,
|
|
tags: options.Outbounds,
|
|
link: options.URL,
|
|
interval: time.Duration(options.Interval),
|
|
tolerance: options.Tolerance,
|
|
idleTimeout: time.Duration(options.IdleTimeout),
|
|
interruptExternalConnections: options.InterruptExistConnections,
|
|
balancer: balancer,
|
|
// nil/absent means true (self-check on) — the documented default, so a
|
|
// config written before the flag existed behaves exactly as it always has.
|
|
selfCheckDisabled: options.SelfCheck != nil && !*options.SelfCheck,
|
|
// Absent from the registry (plain sing-box, tests) yields nil, which
|
|
// means "no gate" — every scheduled probe proceeds, as before.
|
|
probeGate: service.FromContext[urltest.ProbeGate](ctx),
|
|
}
|
|
if len(outbound.tags) == 0 {
|
|
return nil, E.New("missing tags")
|
|
}
|
|
return outbound, nil
|
|
}
|
|
|
|
func (s *URLTest) Start() error {
|
|
outbounds := make([]adapter.Outbound, 0, len(s.tags))
|
|
for i, tag := range s.tags {
|
|
detour, loaded := s.outbound.Outbound(tag)
|
|
if !loaded {
|
|
return E.New("outbound ", i, " not found: ", tag)
|
|
}
|
|
outbounds = append(outbounds, detour)
|
|
}
|
|
group, err := NewURLTestGroup(s.ctx, s.outbound, s.logger, outbounds, s.link, s.interval, s.tolerance, s.idleTimeout, s.interruptExternalConnections)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
group.balancer = s.balancer // lx: SPEC 019 v2 — health-check drives the pool through it
|
|
// lx: health board §5.C — carry the stand-down flag onto the group the same
|
|
// way the balancer travels: set after construction, immutable from then on.
|
|
group.selfCheckDisabled = s.selfCheckDisabled
|
|
// The gate and the tag to ask it about travel together: the gate answers
|
|
// per-outbound, and the group is the thing whose schedule is being gated.
|
|
group.probeGate = s.probeGate
|
|
group.tag = s.Tag()
|
|
if s.balancer != nil {
|
|
// lx: health board §5.B — slot liveness reads through the board verdict, so a
|
|
// death recorded by any prober or a failed dial takes effect on the next pick,
|
|
// not the next health-check tick.
|
|
s.balancer.verdict = group.slotVerdict
|
|
// lx: SPEC 020 — a pool rebuild changes the active routing tree; invalidate
|
|
// the router's reachable cache. ctx captured here has the invalidator.
|
|
ctx := s.ctx
|
|
s.balancer.onChange = func() { invalidateReachability(ctx) }
|
|
}
|
|
s.group = group
|
|
return nil
|
|
}
|
|
|
|
func (s *URLTest) PostStart() error {
|
|
s.group.PostStart()
|
|
return nil
|
|
}
|
|
|
|
func (s *URLTest) Close() error {
|
|
return common.Close(
|
|
common.PtrOrNil(s.group),
|
|
)
|
|
}
|
|
|
|
func (s *URLTest) Now() string {
|
|
// lx: SPEC 019 — balanced modes have no single "current" node; report the last picked tag.
|
|
if s.balancer != nil {
|
|
return s.group.lastSelected.Load()
|
|
}
|
|
// One load, so the two halves reported here are the SAME decision.
|
|
selected := s.group.selected.Load()
|
|
if selected.tcp != nil {
|
|
return selected.tcp.Tag()
|
|
} else if selected.udp != nil {
|
|
return selected.udp.Tag()
|
|
}
|
|
// lx: SPEC 019 — cold start: before the first URL-test the pair is empty but
|
|
// traffic already flows via the Select() fallback (outbounds[0] when no history yet).
|
|
// Mirror exactly what the next DialContext would pick, so the UI shows the real node
|
|
// instead of blank. Select() is the same source of truth DialContext uses.
|
|
if outbound, _ := s.group.Select(N.NetworkTCP); outbound != nil {
|
|
return outbound.Tag()
|
|
}
|
|
if outbound, _ := s.group.Select(N.NetworkUDP); outbound != nil {
|
|
return outbound.Tag()
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func (s *URLTest) All() []string {
|
|
return s.tags
|
|
}
|
|
|
|
// PoolSlot is one entry of the round_robin rotation pool. lx: SPEC 019 v2.
|
|
type PoolSlot struct {
|
|
Slot int
|
|
Tag string
|
|
Delay uint16 // ms; 0 = not measured / dead. A living node is clamped to >= 1 (see Pool).
|
|
}
|
|
|
|
// Pool returns the current rotation pool (one entry per slot) for round_robin groups. For
|
|
// least_test (nil balancer) it returns nil — "this group has no pool". Delay is reported only
|
|
// for slots with a fresh-alive board verdict, clamped 0->1, so 0 in the output unambiguously
|
|
// means dead/untested (failures now persist in history — an entry alone no longer means alive).
|
|
// lx: SPEC 019 v2 (exposed to clients via the GetPool RPC); health board §5.B.
|
|
func (s *URLTest) Pool() []PoolSlot {
|
|
if s.balancer == nil || s.group == nil {
|
|
return nil
|
|
}
|
|
tags := s.balancer.poolTags()
|
|
slots := make([]PoolSlot, len(tags))
|
|
for i, tag := range tags {
|
|
var delay uint16
|
|
// History is keyed by RealTag(detour) (a nested-group member is tested and
|
|
// stored under its live leaf, not the group tag); read it the same way, or
|
|
// the slot's delay is always 0/dead for group members (SPEC 022 #5). Fall
|
|
// back to the raw slot tag if the outbound can't be resolved.
|
|
historyTag := tag
|
|
if node, loaded := s.outbound.Outbound(tag); loaded {
|
|
historyTag = RealTag(node)
|
|
}
|
|
if s.group.history.Verdict(historyTag, s.group.healthTTL()) == urltest.VerdictAlive {
|
|
if history := s.group.history.LoadURLTestHistory(historyTag); history != nil {
|
|
delay = history.Delay
|
|
if delay == 0 {
|
|
delay = 1 // live sub-ms node: never report 0 (0 is reserved for dead/untested)
|
|
}
|
|
}
|
|
}
|
|
slots[i] = PoolSlot{Slot: i, Tag: tag, Delay: delay}
|
|
}
|
|
return slots
|
|
}
|
|
|
|
func (s *URLTest) URLTest(ctx context.Context) (map[string]uint16, error) {
|
|
return s.group.URLTest(ctx)
|
|
}
|
|
|
|
func (s *URLTest) CheckOutbounds() {
|
|
s.group.CheckOutbounds(true)
|
|
}
|
|
|
|
// selectBalanced picks an outbound per-connection in round_robin mode from the balancer's
|
|
// fixed-size pool. lx: SPEC 019 v2. fallback (selectExcluding) covers the cold-start window
|
|
// before the first health-check fills the pool and the all-slots-dead state. exclude carries
|
|
// the members already tried by this connection's dial-retry loop (health board §5.B).
|
|
// Returns nil only when nothing usable.
|
|
func (s *URLTest) selectBalanced(ctx context.Context, network string, destination M.Socksaddr, exclude map[string]bool) adapter.Outbound {
|
|
fallback, _ := s.group.selectExcluding(network, exclude)
|
|
selected := s.balancer.pick(ctx, destination, fallback, func(tag string) adapter.Outbound {
|
|
node, _ := s.outbound.Outbound(tag)
|
|
if node != nil && !common.Contains(node.Network(), network) {
|
|
return nil
|
|
}
|
|
return node
|
|
})
|
|
if selected != nil {
|
|
s.group.lastSelected.Store(selected.Tag())
|
|
}
|
|
return selected
|
|
}
|
|
|
|
func (s *URLTest) DialContext(ctx context.Context, network string, destination M.Socksaddr) (net.Conn, error) {
|
|
s.group.Touch()
|
|
switch N.NetworkName(network) {
|
|
case N.NetworkTCP, N.NetworkUDP:
|
|
default:
|
|
return nil, E.Extend(N.ErrUnknownNetwork, network)
|
|
}
|
|
// lx: health board §5.B — a failed dial no longer fails the user's connection outright:
|
|
// the member is marked dead on the board and the dial is retried through the next
|
|
// candidate (at most dialAttemptsMax members, within the context deadline), for every
|
|
// mode. The old behaviour deleted the history entry in least_test (making a dead member
|
|
// indistinguishable from an untested one) and did nothing at all in balanced modes
|
|
// (plan §2 Д1/Д5).
|
|
var tried map[string]bool
|
|
var lastErr error
|
|
for range dialAttemptsMax {
|
|
outbound := s.dialSelect(ctx, network, destination, tried)
|
|
if outbound == nil {
|
|
break
|
|
}
|
|
conn, err := outbound.DialContext(ctx, network, destination)
|
|
if err == nil {
|
|
return s.group.interruptGroup.NewConn(conn, interrupt.IsExternalConnectionFromContext(ctx)), nil
|
|
}
|
|
s.logger.ErrorContext(ctx, err)
|
|
if tried == nil {
|
|
tried = make(map[string]bool, dialAttemptsMax)
|
|
}
|
|
s.markDialFailure(outbound, tried, err)
|
|
lastErr = err
|
|
if ctx.Err() != nil {
|
|
break
|
|
}
|
|
}
|
|
if lastErr != nil {
|
|
return nil, lastErr
|
|
}
|
|
return nil, E.New("missing supported outbound")
|
|
}
|
|
|
|
func (s *URLTest) ListenPacket(ctx context.Context, destination M.Socksaddr) (net.PacketConn, error) {
|
|
s.group.Touch()
|
|
// lx: health board §5.B — same mark-fail + re-pick retry as DialContext, but only for
|
|
// the ListenPacket call itself: once the packet conn is returned no retry is possible
|
|
// (the UDP session is bound to its member from the first send).
|
|
var tried map[string]bool
|
|
var lastErr error
|
|
for range dialAttemptsMax {
|
|
outbound := s.dialSelect(ctx, N.NetworkUDP, destination, tried)
|
|
if outbound == nil {
|
|
break
|
|
}
|
|
conn, err := outbound.ListenPacket(ctx, destination)
|
|
if err == nil {
|
|
return s.group.interruptGroup.NewPacketConn(conn, interrupt.IsExternalConnectionFromContext(ctx)), nil
|
|
}
|
|
s.logger.ErrorContext(ctx, err)
|
|
if tried == nil {
|
|
tried = make(map[string]bool, dialAttemptsMax)
|
|
}
|
|
s.markDialFailure(outbound, tried, err)
|
|
lastErr = err
|
|
if ctx.Err() != nil {
|
|
break
|
|
}
|
|
}
|
|
if lastErr != nil {
|
|
return nil, lastErr
|
|
}
|
|
return nil, E.New("missing supported outbound")
|
|
}
|
|
|
|
func (s *URLTest) NewConnection(ctx context.Context, conn net.Conn, metadata adapter.InboundContext, onClose N.CloseHandlerFunc) {
|
|
ctx = interrupt.ContextWithIsExternalConnection(ctx)
|
|
s.connection.NewConnection(ctx, s, conn, metadata, onClose)
|
|
}
|
|
|
|
func (s *URLTest) NewPacketConnection(ctx context.Context, conn N.PacketConn, metadata adapter.InboundContext, onClose N.CloseHandlerFunc) {
|
|
ctx = interrupt.ContextWithIsExternalConnection(ctx)
|
|
s.connection.NewPacketConnection(ctx, s, conn, metadata, onClose)
|
|
}
|
|
|
|
type URLTestGroup struct {
|
|
ctx context.Context
|
|
outbound adapter.OutboundManager
|
|
pause pause.Manager
|
|
pauseCallback *list.Element[pause.Callback]
|
|
logger log.Logger
|
|
outbounds []adapter.Outbound
|
|
link string
|
|
interval time.Duration
|
|
tolerance uint16
|
|
idleTimeout time.Duration
|
|
history *urltest.HistoryStorage
|
|
checking atomic.Bool
|
|
// selected is the least_test cache: the member this group currently prefers, per
|
|
// network. It is written by the probing goroutine and read on EVERY dial through
|
|
// the group (selectExcluding / dialSelect) and by the panel (Now), so it is an
|
|
// atomic value rather than two plain fields — the same thing Selector does one file
|
|
// over (selector.go, common.TypedValue[adapter.Outbound]). An interface field is two
|
|
// words; a torn read of one hands a dial a type descriptor with the wrong data
|
|
// pointer, which is not a wrong node but a corrupt one.
|
|
//
|
|
// The TCP and UDP halves live in ONE value on purpose. They are decided together, by
|
|
// one pass over one board reading, and publishing them separately let a reader pick
|
|
// up the new TCP choice against the previous UDP choice — the group's hysteresis
|
|
// silently applied to a decision that was never made.
|
|
selected common.TypedValue[selectedPair]
|
|
interruptGroup *interrupt.Group
|
|
interruptExternalConnections bool
|
|
access sync.Mutex
|
|
ticker *time.Ticker
|
|
close chan struct{}
|
|
// started is read by Touch on every dial, outside g.access, and written by
|
|
// PostStart under it — an atomic because that is what it always was in effect.
|
|
started atomic.Bool
|
|
// closed latches in Close and is what makes Close FINAL. Guarded by access.
|
|
//
|
|
// It exists because "has a ticker" is not the same question as "is shut down", and
|
|
// Close used to ask the first one: with no ticker armed it returned before closing
|
|
// g.close, leaving the group indistinguishable from a running one. A Touch arriving
|
|
// afterwards — an outbound snapshot taken before an Apply is still dialable for up
|
|
// to two minutes, see shater/engine/grouptest.go — then armed a fresh ticker whose
|
|
// loopCheck waits on a channel nobody will ever close, in a box whose context is
|
|
// already cancelled. Every tick of it fails instantly and files a "dead" verdict on
|
|
// the SHARED health board that the live generation selects nodes from. One retired
|
|
// group can go on declaring the whole node set dead for the uptime of the daemon.
|
|
closed bool
|
|
lastActive common.TypedValue[time.Time]
|
|
lastSelected common.TypedValue[string] // lx: SPEC 019 — Now() in balanced modes
|
|
balancer *balancer // lx: SPEC 019 v2 — round_robin pool; nil for least_test
|
|
// lx: health board §5.C — mirrors URLTest.selfCheckDisabled (set by Start,
|
|
// immutable afterwards, zero value = probing on). Guards ONLY the group's
|
|
// own schedule: the PostStart warm-up sweep and the Touch ticker. An
|
|
// explicit CheckOutbounds/URLTest call is untouched — the flag stands down
|
|
// the schedule, not the capability.
|
|
selfCheckDisabled bool
|
|
// lx: health board §5.C — the runtime gate and the tag it is asked about.
|
|
// Both mirror URLTest's fields (set by Start, immutable afterwards); nil
|
|
// gate or empty tag means every scheduled check proceeds. Consulted only
|
|
// through selfCheckAllowed, and only on the SCHEDULE.
|
|
probeGate urltest.ProbeGate
|
|
tag string
|
|
}
|
|
|
|
// selectedPair is one published least_test decision: the member chosen for TCP and the
|
|
// member chosen for UDP, as of the same probing round. Either half may be nil (nothing
|
|
// picked yet for that network).
|
|
type selectedPair struct {
|
|
tcp adapter.Outbound
|
|
udp adapter.Outbound
|
|
}
|
|
|
|
// selectedFor returns the cached choice for one network (nil when there is none, or when
|
|
// network is neither TCP nor UDP — the caller then falls through to a fresh selection).
|
|
func (g *URLTestGroup) selectedFor(network string) adapter.Outbound {
|
|
pair := g.selected.Load()
|
|
switch network {
|
|
case N.NetworkTCP:
|
|
return pair.tcp
|
|
case N.NetworkUDP:
|
|
return pair.udp
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// setSelected publishes a decision. It is the ONLY writer of g.selected, and it writes
|
|
// the pair whole — see the field comment for why the two halves may not be split.
|
|
func (g *URLTestGroup) setSelected(tcp, udp adapter.Outbound) {
|
|
g.selected.Store(selectedPair{tcp: tcp, udp: udp})
|
|
}
|
|
|
|
// selfCheckAllowed reports whether the group's OWN probing schedule may dial
|
|
// right now. lx: health board §5.C.
|
|
//
|
|
// Two independent refusals, in the order they can be answered cheapest first:
|
|
//
|
|
// selfCheckDisabled — the config says no rule reaches this group. Fixed for
|
|
// the life of the box; see standDownUnusedSelfCheck.
|
|
// probeGate — the world says this group cannot be reached right now,
|
|
// typically a chain hop sitting behind a dead hop. Asked
|
|
// fresh EVERY time, which is the entire mechanism by which
|
|
// a recovered hop resumes probing: there is no state here
|
|
// to reset, so there is none to get stuck.
|
|
//
|
|
// Neither refusal touches an explicit CheckOutbounds/URLTest — a deliberate
|
|
// request is never a scheduled one.
|
|
func (g *URLTestGroup) selfCheckAllowed() bool {
|
|
if g.selfCheckDisabled {
|
|
return false
|
|
}
|
|
if g.probeGate == nil || g.tag == "" {
|
|
return true
|
|
}
|
|
return g.probeGate.ProbeAllowed(g.tag)
|
|
}
|
|
|
|
// scheduledCheck is one firing of the group's own schedule — the warm-up sweep
|
|
// and every ticker tick go through here, and nothing else does. Having exactly
|
|
// one gated entry point is what keeps the two callers from drifting apart, and
|
|
// it is the seam the tests drive to assert that a gated group makes no dial
|
|
// attempt at all.
|
|
func (g *URLTestGroup) scheduledCheck() {
|
|
if !g.selfCheckAllowed() {
|
|
return
|
|
}
|
|
g.CheckOutbounds(false)
|
|
}
|
|
|
|
// keepWarm reports whether this group must keep measuring with no traffic
|
|
// flowing through it. lx: health board §5.C — see urltest.ProbeGate.ProbeWhenIdle.
|
|
//
|
|
// The default is NO, in every direction: no gate, no tag, or a group whose
|
|
// self-check is stood down anyway. Only a gate that positively says "the routing
|
|
// config reaches this group" turns the idle timeout off, so plain sing-box and
|
|
// every hand-built group keep the lifecycle they have always had.
|
|
func (g *URLTestGroup) keepWarm() bool {
|
|
if g.selfCheckDisabled || g.probeGate == nil || g.tag == "" {
|
|
return false
|
|
}
|
|
return g.probeGate.ProbeWhenIdle(g.tag)
|
|
}
|
|
|
|
// startTickerLocked arms the group's own probing ticker. g.access MUST be held
|
|
// and g.ticker MUST be nil. Extracted so PostStart and Touch arm it identically
|
|
// — two ways in, one construction, no chance of one of them forgetting the pause
|
|
// registration.
|
|
func (g *URLTestGroup) startTickerLocked() {
|
|
ticker := time.NewTicker(g.interval)
|
|
g.ticker = ticker
|
|
g.pauseCallback = pause.RegisterTicker(g.pause, ticker, g.interval, nil)
|
|
go g.loopCheck(ticker, g.close)
|
|
}
|
|
|
|
func NewURLTestGroup(ctx context.Context, outboundManager adapter.OutboundManager, logger log.Logger, outbounds []adapter.Outbound, link string, interval time.Duration, tolerance uint16, idleTimeout time.Duration, interruptExternalConnections bool) (*URLTestGroup, error) {
|
|
if interval == 0 {
|
|
interval = C.DefaultURLTestInterval
|
|
}
|
|
if tolerance == 0 {
|
|
tolerance = 50
|
|
}
|
|
if idleTimeout == 0 {
|
|
idleTimeout = C.DefaultURLTestIdleTimeout
|
|
}
|
|
if interval > idleTimeout {
|
|
return nil, E.New("interval must be less or equal than idle_timeout")
|
|
}
|
|
history := service.PtrFromContext[urltest.HistoryStorage](ctx)
|
|
if history == nil {
|
|
return nil, E.New("missing URL test history storage")
|
|
}
|
|
return &URLTestGroup{
|
|
ctx: ctx,
|
|
outbound: outboundManager,
|
|
logger: logger,
|
|
outbounds: outbounds,
|
|
link: link,
|
|
interval: interval,
|
|
tolerance: tolerance,
|
|
idleTimeout: idleTimeout,
|
|
history: history,
|
|
close: make(chan struct{}),
|
|
pause: service.FromContext[pause.Manager](ctx),
|
|
interruptGroup: interrupt.NewGroup(),
|
|
interruptExternalConnections: interruptExternalConnections,
|
|
}, nil
|
|
}
|
|
|
|
func (g *URLTestGroup) PostStart() {
|
|
g.access.Lock()
|
|
defer g.access.Unlock()
|
|
if g.closed {
|
|
return
|
|
}
|
|
g.started.Store(true)
|
|
g.lastActive.Store(time.Now())
|
|
// lx: SPEC 019 v2 — seed the pool so round_robin can route from the first connection,
|
|
// before the first health-check completes (history-warm nodes first, else config order).
|
|
// The seed only READS the board, so it runs even with the self-check stood down.
|
|
g.seedPool()
|
|
// lx: health board §5.C — the warm-up sweep is the first half of the group's
|
|
// own probing schedule, and it fires for EVERY group at box start, including
|
|
// groups no routing rule reaches. For those, the sweep dials every member
|
|
// directly from the router — a path nothing uses — and records the outcome
|
|
// under the members' base tags, forging the board reading the observatory
|
|
// exists to keep honest. A stood-down group therefore skips it entirely; the
|
|
// observatory (or nothing, for a truly unused group) is what measures its
|
|
// members.
|
|
//
|
|
// The same call is now also where a chain hop behind a DEAD hop declines to
|
|
// sweep: every member of such a group dials through the broken hop, so the
|
|
// sweep would measure that hop once per member and file the result against
|
|
// this one. selfCheckAllowed keeps both refusals in one place.
|
|
go g.scheduledCheck()
|
|
// A group the routing config REACHES keeps measuring whether or not anybody
|
|
// dials it, so its ticker is armed here instead of waiting for a Touch that
|
|
// may never come. Without this, a used group with no traffic gets this one
|
|
// warm-up sweep and then nothing: its members age past the verdict TTL and
|
|
// the panel reports "untested" about a rule that is in force, while the first
|
|
// real request pays a cold probe. Nothing else would fill the gap — the
|
|
// observatory stands off a urltest group's members entirely (probeplan.go
|
|
// SelfChecked), which is the whole point of one dialler per target.
|
|
//
|
|
// lastActive was stored a moment ago, so loopCheck's opening "idle longer
|
|
// than the interval" check does not fire and this cannot double up with the
|
|
// sweep above.
|
|
if g.keepWarm() && g.ticker == nil {
|
|
g.startTickerLocked()
|
|
}
|
|
}
|
|
|
|
func (g *URLTestGroup) Touch() {
|
|
if !g.started.Load() {
|
|
return
|
|
}
|
|
// lx: health board §5.C — Touch's only job is to keep the group's OWN
|
|
// probing ticker alive while traffic flows. With the self-check stood down
|
|
// there is deliberately no ticker to start or feed: the observatory owns the
|
|
// schedule, and a stray dial through an unused group (a stale rule cache, a
|
|
// manual pin) must not arm 30 minutes of direct probing under the members'
|
|
// base tags. Checked before the lock because the flag is immutable after
|
|
// Start, exactly like the started fast-path above.
|
|
//
|
|
// The runtime gate is deliberately NOT consulted here. Touch only arms the
|
|
// ticker; refusing to arm it would mean a hop that recovers has no ticker
|
|
// left to notice — the block would outlive the failure, which is the one
|
|
// outcome this must never have. The ticker runs and each tick re-asks the
|
|
// gate (loopCheck -> scheduledCheck), so a blocked hop costs a predicate
|
|
// call per interval and resumes the moment the hop in front answers.
|
|
if g.selfCheckDisabled {
|
|
return
|
|
}
|
|
g.access.Lock()
|
|
defer g.access.Unlock()
|
|
// A closed group arms nothing. Touch is reachable long after Close — a caller
|
|
// holding an outbound from a snapshot taken before an Apply keeps dialling it (up
|
|
// to the 120s budget of shater/engine/grouptest.go) — and the ticker it would arm
|
|
// has no way left to stop: see the `closed` field for what that costs.
|
|
if g.closed {
|
|
return
|
|
}
|
|
if g.ticker != nil {
|
|
g.lastActive.Store(time.Now())
|
|
return
|
|
}
|
|
g.startTickerLocked()
|
|
}
|
|
|
|
// Close shuts the group down for good. It is idempotent, and it is FINAL: no later Touch
|
|
// can bring the probing schedule back.
|
|
//
|
|
// It used to return early when no ticker happened to be armed, without ever closing
|
|
// g.close — so a group that was closed while idle stayed, from the point of view of every
|
|
// other method, a perfectly live group. That is the whole defect: the close channel is the
|
|
// only way a loopCheck goroutine ever exits (its idle-timeout escape does not fire for a
|
|
// group the routing config reaches, keepWarm), so a ticker armed after such a Close is
|
|
// immortal, and every one of its ticks writes a failure to the shared health board on
|
|
// behalf of a box that no longer exists.
|
|
func (g *URLTestGroup) Close() error {
|
|
g.access.Lock()
|
|
defer g.access.Unlock()
|
|
if g.closed {
|
|
return nil
|
|
}
|
|
g.closed = true
|
|
// Unconditionally, BEFORE looking at the ticker: this is the signal every loopCheck
|
|
// waits on, including any that a Touch armed after the last one was retired by the
|
|
// idle timeout.
|
|
if g.close != nil {
|
|
close(g.close)
|
|
}
|
|
if g.ticker != nil {
|
|
g.ticker.Stop()
|
|
g.ticker = nil
|
|
g.pause.UnregisterCallback(g.pauseCallback)
|
|
g.pauseCallback = nil
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (g *URLTestGroup) Select(network string) (adapter.Outbound, bool) {
|
|
// lx: health board §5.B — selection reads the board verdict instead of "has a history
|
|
// entry": fresh-alive members ranked by delay first, untested members in config order
|
|
// second, first-by-config last — a dead member is never picked while a live or untested
|
|
// one exists. The logic lives in selectExcluding (urltest_health_lx.go) so the
|
|
// dial-retry path can re-run it minus the members that just failed.
|
|
return g.selectExcluding(network, nil)
|
|
}
|
|
|
|
// loopCheck is the group's own schedule. lx: health board §5.C — every probe it
|
|
// fires goes through scheduledCheck, so a stood-down or currently-unreachable
|
|
// group ticks without dialling. The ticker's LIFECYCLE (the idle timeout below)
|
|
// is deliberately left alone: a gated group keeps its ticker exactly as long as
|
|
// an ungated one would, because the ticker is what will notice the recovery.
|
|
func (g *URLTestGroup) loopCheck(ticker *time.Ticker, closeChan <-chan struct{}) {
|
|
if time.Since(g.lastActive.Load()) > g.interval {
|
|
g.lastActive.Store(time.Now())
|
|
g.scheduledCheck()
|
|
}
|
|
for {
|
|
select {
|
|
case <-closeChan:
|
|
return
|
|
case <-ticker.C:
|
|
}
|
|
// The idle timeout retires the ticker of a group nobody is dialling —
|
|
// unless the routing config reaches it, in which case its health is a
|
|
// live question whether or not traffic is flowing and the ticker must
|
|
// outlive the silence. Asked here rather than remembered from PostStart
|
|
// so it tracks the running config, and asked OUTSIDE g.access because the
|
|
// answer comes from the engine, which has locks of its own.
|
|
if !g.keepWarm() && time.Since(g.lastActive.Load()) > g.idleTimeout {
|
|
g.access.Lock()
|
|
if g.ticker == ticker {
|
|
g.ticker.Stop()
|
|
g.ticker = nil
|
|
g.pause.UnregisterCallback(g.pauseCallback)
|
|
g.pauseCallback = nil
|
|
}
|
|
g.access.Unlock()
|
|
return
|
|
}
|
|
g.scheduledCheck()
|
|
}
|
|
}
|
|
|
|
// lx: health board §5.B/§5.C — CheckOutbounds is the FAST loop of the pair: an
|
|
// ACTIVE group probes its own members on its own ticker (interval = the global
|
|
// probe interval; failover keeps its 30s default), and a failed probe writes
|
|
// MarkFailed to the shared health board (common/urltest) instead of deleting
|
|
// the entry. The shater observatory is the complementary BACKGROUND loop: it
|
|
// probes what no active group measures (idle groups' members, selector
|
|
// candidates, chain exits), and its freshness gate skips any tag this loop
|
|
// keeps current — no duplicate probes, and this loop's cadence is never
|
|
// suppressed. Both loops write to the SAME board, so selection and the panel
|
|
// see one source of truth no matter which prober found the death.
|
|
func (g *URLTestGroup) CheckOutbounds(force bool) {
|
|
_, _ = g.urlTest(g.ctx, force)
|
|
}
|
|
|
|
func (g *URLTestGroup) URLTest(ctx context.Context) (map[string]uint16, error) {
|
|
return g.urlTest(ctx, false)
|
|
}
|
|
|
|
func (g *URLTestGroup) urlTest(ctx context.Context, force bool) (map[string]uint16, error) {
|
|
result := make(map[string]uint16)
|
|
if g.checking.Swap(true) {
|
|
return result, nil
|
|
}
|
|
defer g.checking.Store(false)
|
|
// lx: SPEC 019 v2 — round_robin uses a lazy, pool-bounded health-check that tests no more
|
|
// nodes than needed (unless force, e.g. a manual URLTest, which always tests everything).
|
|
if g.balancer != nil && !force {
|
|
return g.balancePool(ctx), nil
|
|
}
|
|
result = g.testNodes(ctx, g.outbounds, force)
|
|
if g.balancer != nil {
|
|
// force path (manual URLTest tested all nodes): rebuild the pool from fresh results.
|
|
g.rebuildPool()
|
|
} else {
|
|
g.performUpdateCheck()
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
// testNodes runs the URL test over the given outbounds (skipping fresh history unless force),
|
|
// stores successes and marks failures on the health board, and returns tag->delay for the live
|
|
// ones. lx: shared by least_test and the round_robin force path; this is the probe
|
|
// primitive of the fast active-group loop (see CheckOutbounds for the two-loop
|
|
// contract with the shater observatory).
|
|
func (g *URLTestGroup) testNodes(ctx context.Context, outbounds []adapter.Outbound, force bool) map[string]uint16 {
|
|
result := make(map[string]uint16)
|
|
b, _ := batch.New(ctx, batch.WithConcurrencyNum[any](10))
|
|
checked := make(map[string]bool)
|
|
var resultAccess sync.Mutex
|
|
for _, detour := range outbounds {
|
|
tag := detour.Tag()
|
|
realTag := RealTag(detour)
|
|
if checked[realTag] {
|
|
continue
|
|
}
|
|
history := g.history.LoadURLTestHistory(realTag)
|
|
if !force && history != nil && time.Since(history.LastOK) < g.interval { // lx: health board §5.A — Time renamed to LastOK
|
|
continue
|
|
}
|
|
checked[realTag] = true
|
|
p, loaded := g.outbound.Outbound(realTag)
|
|
if !loaded {
|
|
continue
|
|
}
|
|
b.Go(realTag, func() (any, error) {
|
|
testCtx, cancel := context.WithTimeout(g.ctx, C.TCPTimeout)
|
|
defer cancel()
|
|
t, err := urltest.URLTest(testCtx, g.link, p)
|
|
if err != nil {
|
|
g.logger.Debug("outbound ", tag, " unavailable: ", err)
|
|
// lx: health board §5.A — a failed check marks the entry instead of deleting
|
|
// it (deletion stays reserved for members removed from the configuration).
|
|
g.markFailedLogged(realTag, tag, "probe", err)
|
|
} else {
|
|
g.logger.Debug("outbound ", tag, " available: ", t, "ms")
|
|
// lx: health board — log the dead -> alive flip before the success overwrites it.
|
|
if g.history.Verdict(realTag, g.healthTTL()) == urltest.VerdictDead {
|
|
g.logger.Info("outbound ", tag, " flipped dead -> alive (probe: ", t, "ms)")
|
|
}
|
|
g.history.StoreURLTestHistory(realTag, &adapter.URLTestHistory{
|
|
LastOK: time.Now(), // lx: health board §5.A — Time renamed to LastOK
|
|
Delay: t,
|
|
})
|
|
resultAccess.Lock()
|
|
result[tag] = t
|
|
resultAccess.Unlock()
|
|
}
|
|
return nil, nil
|
|
})
|
|
}
|
|
b.Wait()
|
|
return result
|
|
}
|
|
|
|
// performUpdateCheck re-ranks the members after a probing round and publishes the result.
|
|
// It is the only writer of g.selected: it reads the current pair ONCE, decides both
|
|
// networks against that one snapshot, and stores the outcome as a single value, so no
|
|
// reader can ever observe a half-applied decision. Callers are serialised by g.checking
|
|
// (urlTest), which is what makes the read-decide-store sequence safe without a lock.
|
|
func (g *URLTestGroup) performUpdateCheck() {
|
|
current := g.selected.Load()
|
|
next := current
|
|
var updated bool
|
|
if outbound, exists := g.Select(N.NetworkTCP); outbound != nil && (current.tcp == nil || (exists && outbound != current.tcp)) {
|
|
if current.tcp != nil {
|
|
updated = true
|
|
}
|
|
next.tcp = outbound
|
|
}
|
|
if outbound, exists := g.Select(N.NetworkUDP); outbound != nil && (current.udp == nil || (exists && outbound != current.udp)) {
|
|
if current.udp != nil {
|
|
updated = true
|
|
}
|
|
next.udp = outbound
|
|
}
|
|
if next != current {
|
|
g.setSelected(next.tcp, next.udp)
|
|
}
|
|
if updated {
|
|
g.interruptGroup.Interrupt(g.interruptExternalConnections)
|
|
invalidateReachability(g.ctx) // lx: SPEC 020 — legacy auto-switch changed the active node
|
|
}
|
|
}
|
|
|
|
// --- round_robin pool health-check (lx: SPEC 019 v2) --------------------------------
|
|
|
|
// poolSize is the effective pool size for the current node set: min(configured, available).
|
|
func (g *URLTestGroup) poolSize() int {
|
|
size := g.balancer.poolSize
|
|
if size > len(g.outbounds) {
|
|
size = len(g.outbounds)
|
|
}
|
|
return size
|
|
}
|
|
|
|
// balancePool is the per-interval lazy health-check for round_robin. It tests no more nodes
|
|
// than needed to keep the pool full of live nodes, then applies the new slot occupancy.
|
|
// Returns tag->delay for every node it tested live (for the URLTest map / UI).
|
|
func (g *URLTestGroup) balancePool(ctx context.Context) map[string]uint16 {
|
|
size := g.poolSize()
|
|
if size == 0 {
|
|
return map[string]uint16{}
|
|
}
|
|
if g.balancer.priority {
|
|
return g.balancePoolPriority(ctx, size)
|
|
}
|
|
if g.balancer.poolTolerance > 0 {
|
|
return g.balancePoolTolerant(ctx, size)
|
|
}
|
|
return g.balancePoolFirstLive(ctx, size)
|
|
}
|
|
|
|
// balancePoolPriority (priority balancer, strategy=failover) walks the group members in CONFIG
|
|
// ORDER from the top, testing one at a time, until `size` live nodes are found — then makes those
|
|
// the pool via planPriorityPool. Two properties fall out of testing from the top every tick:
|
|
//
|
|
// - fail-back: a higher-priority node that has come back to life is re-probed (it is above the
|
|
// current occupant in config order, so the walk reaches it first) and re-takes its slot on
|
|
// this tick. This is what balancePoolFirstLive cannot do — that path re-tests the CURRENT
|
|
// pool first and only probes others to fill a hole, so a revived #1 is never reconsidered
|
|
// while #2 stays live.
|
|
// - bounded cost: the walk STOPS at the `size`-th live node, so nodes BELOW it are never
|
|
// probed. For failover (size 1) that is one probe per tick when the top node is up, and one
|
|
// extra probe per dead node above the first live one while the top of the list is down.
|
|
func (g *URLTestGroup) balancePoolPriority(ctx context.Context, size int) map[string]uint16 {
|
|
result := make(map[string]uint16)
|
|
live := make(map[string]bool)
|
|
configOrder := make([]string, 0, len(g.outbounds))
|
|
for _, detour := range g.outbounds {
|
|
configOrder = append(configOrder, detour.Tag())
|
|
}
|
|
liveCount := 0
|
|
for _, detour := range g.outbounds {
|
|
if liveCount >= size {
|
|
break // enough live nodes for the pool; do not probe those below them.
|
|
}
|
|
tag := detour.Tag()
|
|
tested := g.testNodes(ctx, []adapter.Outbound{detour}, true)
|
|
if delay, ok := tested[tag]; ok {
|
|
result[tag] = delay
|
|
live[tag] = true
|
|
liveCount++
|
|
}
|
|
}
|
|
g.balancer.setSlots(planPriorityPool(configOrder, live, size), live)
|
|
return result
|
|
}
|
|
|
|
// balancePoolFirstLive (pool_tolerance == 0): re-test the nodes already in the pool, then —
|
|
// only if the pool is short of live nodes — walk the rest in config order, testing until the
|
|
// pool is full again. A dead pool node keeps its slot until a live replacement is found.
|
|
func (g *URLTestGroup) balancePoolFirstLive(ctx context.Context, size int) map[string]uint16 {
|
|
current := g.balancer.poolTags()
|
|
inPool := make(map[string]bool, len(current))
|
|
for _, tag := range current {
|
|
if tag != "" {
|
|
inPool[tag] = true
|
|
}
|
|
}
|
|
// 1. Re-test current pool members; collect which slots went dead.
|
|
poolNodes := g.outboundsByTags(current)
|
|
result := g.testNodes(ctx, poolNodes, true)
|
|
liveTag := func(tag string) bool { _, ok := result[tag]; return ok }
|
|
|
|
// 2. Build the next occupancy IN PLACE: a live member keeps its exact slot index; a dead or
|
|
// empty slot becomes "" (a hole to be refilled). Never compact — shifting a living node
|
|
// across slots would move every sticky key bound to it (the SPEC invariant, see the file
|
|
// header in urltest_balance_lx.go). next is at least `size` long so the pool can grow.
|
|
slotCount := size
|
|
if len(current) > slotCount {
|
|
slotCount = len(current)
|
|
}
|
|
next := make([]string, slotCount)
|
|
for i, tag := range current {
|
|
if tag != "" && liveTag(tag) {
|
|
next[i] = tag
|
|
}
|
|
}
|
|
// emptySlot returns the first hole at/after `from`, or -1 when the pool is full.
|
|
emptySlot := func(from int) int {
|
|
for i := from; i < len(next); i++ {
|
|
if next[i] == "" {
|
|
return i
|
|
}
|
|
}
|
|
return -1
|
|
}
|
|
// 2b. Refill holes (dead/empty slots) by writing replacements INTO the hole's own index:
|
|
// walk non-pool nodes in config order, testing in batches of `size` (a full pool's worth)
|
|
// in parallel, until no holes remain or nodes run out.
|
|
if emptySlot(0) >= 0 {
|
|
candidates := make([]adapter.Outbound, 0, len(g.outbounds))
|
|
for _, detour := range g.outbounds {
|
|
if !inPool[detour.Tag()] {
|
|
candidates = append(candidates, detour)
|
|
}
|
|
}
|
|
fill := 0
|
|
for start := 0; start < len(candidates) && emptySlot(fill) >= 0; start += size {
|
|
end := start + size
|
|
if end > len(candidates) {
|
|
end = len(candidates)
|
|
}
|
|
batch := candidates[start:end]
|
|
tested := g.testNodes(ctx, batch, true)
|
|
// Take live ones in config order (batch is already in config order).
|
|
for _, detour := range batch {
|
|
slot := emptySlot(fill)
|
|
if slot < 0 {
|
|
break
|
|
}
|
|
tag := detour.Tag()
|
|
if delay, ok := tested[tag]; ok {
|
|
next[slot] = tag
|
|
fill = slot + 1
|
|
result[tag] = delay
|
|
}
|
|
}
|
|
}
|
|
}
|
|
// 3. Any hole left (not enough live nodes): put a dead member back in it so the pool never
|
|
// shrinks. A dead occupant keeps the slot it already held when possible; otherwise the
|
|
// remaining dead members fill the leftover holes (order does not matter — all are dead).
|
|
if emptySlot(0) >= 0 {
|
|
// Slots that still hold their original dead occupant: leave them be.
|
|
for i, tag := range current {
|
|
if i < len(next) && next[i] == "" && tag != "" && !liveTag(tag) {
|
|
next[i] = tag
|
|
}
|
|
}
|
|
// Surplus dead members (slots that no longer fit) drop into any leftover hole.
|
|
placed := make(map[string]bool, len(next))
|
|
for _, tag := range next {
|
|
if tag != "" {
|
|
placed[tag] = true
|
|
}
|
|
}
|
|
for _, tag := range current {
|
|
slot := emptySlot(0)
|
|
if slot < 0 {
|
|
break
|
|
}
|
|
if tag != "" && !liveTag(tag) && !placed[tag] {
|
|
next[slot] = tag
|
|
placed[tag] = true
|
|
}
|
|
}
|
|
}
|
|
// result holds exactly the tags that tested live this round (pool re-test + hole fills);
|
|
// setSlots marks every other slot dead so pick() skips it. lx: SPEC 019 v2.
|
|
g.balancer.setSlots(next, liveSet(result))
|
|
return result
|
|
}
|
|
|
|
// liveSet turns a tag->delay result map (only live nodes are present) into the tag->live set
|
|
// setSlots consumes. lx: SPEC 019 v2 — pick() routes only through live slots.
|
|
func liveSet(result map[string]uint16) map[string]bool {
|
|
live := make(map[string]bool, len(result))
|
|
for tag := range result {
|
|
live[tag] = true
|
|
}
|
|
return live
|
|
}
|
|
|
|
// balancePoolTolerant (pool_tolerance > 0): test all nodes, then pick the top-`size` by delay,
|
|
// replacing a pool member only when an outside node beats it by more than the tolerance.
|
|
func (g *URLTestGroup) balancePoolTolerant(ctx context.Context, size int) map[string]uint16 {
|
|
result := g.testNodes(ctx, g.outbounds, true)
|
|
results := make(map[string]candidate, len(g.outbounds))
|
|
for _, detour := range g.outbounds {
|
|
tag := detour.Tag()
|
|
if delay, ok := result[tag]; ok {
|
|
results[tag] = candidate{tag: tag, delay: delay, alive: true}
|
|
} else {
|
|
results[tag] = candidate{tag: tag, alive: false}
|
|
}
|
|
}
|
|
next := planTolerantPool(g.balancer.poolTags(), results, size, g.balancer.poolTolerance)
|
|
g.balancer.setSlots(next, liveSet(result))
|
|
return result
|
|
}
|
|
|
|
// rebuildPool re-derives slot occupancy after a forced full test (manual URLTest) from the
|
|
// history now present. It honours pool_tolerance: with tolerance == 0 it keeps living members
|
|
// in their slots (first-live discipline — a manual test must not reshuffle a stable pool and
|
|
// break sticky bindings); with tolerance > 0 it re-ranks by delay like the steady-state path.
|
|
func (g *URLTestGroup) rebuildPool() {
|
|
size := g.poolSize()
|
|
if size == 0 {
|
|
return
|
|
}
|
|
results := make(map[string]candidate, len(g.outbounds))
|
|
for _, detour := range g.outbounds {
|
|
tag := detour.Tag()
|
|
realTag := RealTag(detour)
|
|
// lx: health board §5.B — alive = fresh board verdict, not entry presence
|
|
// (failures persist in history now).
|
|
if g.history.Verdict(realTag, g.healthTTL()) == urltest.VerdictAlive {
|
|
history := g.history.LoadURLTestHistory(realTag)
|
|
results[tag] = candidate{tag: tag, delay: history.Delay, alive: true}
|
|
} else {
|
|
results[tag] = candidate{tag: tag, alive: false}
|
|
}
|
|
}
|
|
current := g.balancer.poolTags()
|
|
if g.balancer.priority {
|
|
// Priority (failover): config order IS the ranking, so a forced full test rebuilds the
|
|
// pool straight from planPriorityPool — the first `size` live members in config order.
|
|
// This also re-applies fail-back after a manual "Test all nodes".
|
|
live := make(map[string]bool, len(results))
|
|
for tag, c := range results {
|
|
if c.alive {
|
|
live[tag] = true
|
|
}
|
|
}
|
|
configOrder := make([]string, 0, len(g.outbounds))
|
|
for _, detour := range g.outbounds {
|
|
configOrder = append(configOrder, detour.Tag())
|
|
}
|
|
g.balancer.setSlots(planPriorityPool(configOrder, live, size), live)
|
|
return
|
|
}
|
|
if g.balancer.poolTolerance == 0 {
|
|
live := make(map[string]bool, len(results))
|
|
inPool := make(map[string]bool, len(current))
|
|
for _, tag := range current {
|
|
if tag != "" {
|
|
inPool[tag] = true
|
|
}
|
|
}
|
|
for tag, c := range results {
|
|
if c.alive {
|
|
live[tag] = true
|
|
}
|
|
}
|
|
// Fill holes with live non-pool nodes, fastest first (history is warm after the force test).
|
|
fillCandidates := make([]candidate, 0, len(results))
|
|
for _, detour := range g.outbounds {
|
|
tag := detour.Tag()
|
|
if c, ok := results[tag]; ok && c.alive && !inPool[tag] {
|
|
fillCandidates = append(fillCandidates, c)
|
|
}
|
|
}
|
|
sortCandidatesByDelay(fillCandidates)
|
|
fillOrder := make([]string, len(fillCandidates))
|
|
for i, c := range fillCandidates {
|
|
fillOrder[i] = c.tag
|
|
}
|
|
g.balancer.setSlots(planFirstLivePool(current, live, fillOrder, size), live)
|
|
return
|
|
}
|
|
tolerantLive := make(map[string]bool, len(results))
|
|
for tag, c := range results {
|
|
if c.alive {
|
|
tolerantLive[tag] = true
|
|
}
|
|
}
|
|
g.balancer.setSlots(planTolerantPool(current, results, size, g.balancer.poolTolerance), tolerantLive)
|
|
}
|
|
|
|
// seedPool fills the pool before the first health-check: prefer nodes the board holds a
|
|
// fresh-alive verdict for (the process was not unloaded), else the first `size` nodes in
|
|
// config order. lx: SPEC 019 v2; health board §5.B — a persisted failure record must not
|
|
// look like a warm node.
|
|
func (g *URLTestGroup) seedPool() {
|
|
if g.balancer == nil {
|
|
return
|
|
}
|
|
size := g.poolSize()
|
|
if size == 0 {
|
|
return
|
|
}
|
|
// Fresh-alive nodes first (top by delay), then config order to fill.
|
|
withHistory := make([]candidate, 0, len(g.outbounds))
|
|
for _, detour := range g.outbounds {
|
|
realTag := RealTag(detour)
|
|
if g.history.Verdict(realTag, g.healthTTL()) != urltest.VerdictAlive {
|
|
continue
|
|
}
|
|
if history := g.history.LoadURLTestHistory(realTag); history != nil {
|
|
withHistory = append(withHistory, candidate{tag: detour.Tag(), delay: history.Delay, alive: true})
|
|
}
|
|
}
|
|
sortCandidatesByDelay(withHistory)
|
|
next := make([]string, 0, size)
|
|
seen := make(map[string]bool)
|
|
for _, c := range withHistory {
|
|
if len(next) >= size {
|
|
break
|
|
}
|
|
next = append(next, c.tag)
|
|
seen[c.tag] = true
|
|
}
|
|
for _, detour := range g.outbounds {
|
|
if len(next) >= size {
|
|
break
|
|
}
|
|
tag := detour.Tag()
|
|
if !seen[tag] {
|
|
next = append(next, tag)
|
|
seen[tag] = true
|
|
}
|
|
}
|
|
// Seed marks every seeded slot LIVE optimistically: no health-check has run yet, so we have
|
|
// no failure evidence, and marking them dead would blackhole all traffic to the fallback until
|
|
// the first check completes (killing cold-start spread). The first CheckOutbounds (kicked off
|
|
// right after seedPool in PostStart) corrects any that are actually down within one interval.
|
|
live := make(map[string]bool, len(next))
|
|
for _, tag := range next {
|
|
live[tag] = true
|
|
}
|
|
g.balancer.setSlots(next, live)
|
|
}
|
|
|
|
// outboundsByTags resolves slot tags back to live outbound objects (skipping empties/unknowns).
|
|
func (g *URLTestGroup) outboundsByTags(tags []string) []adapter.Outbound {
|
|
out := make([]adapter.Outbound, 0, len(tags))
|
|
for _, tag := range tags {
|
|
if tag == "" {
|
|
continue
|
|
}
|
|
if node, ok := g.outbound.Outbound(tag); ok {
|
|
out = append(out, node)
|
|
}
|
|
}
|
|
return out
|
|
}
|