Files
shater/protocol/group/urltest.go
T
omar a57717dabb
release / aarch64_cortex-a53 (push) Successful in 3m43s
release / x86_64 (push) Successful in 3m30s
release / apk aarch64_cortex-a53 (push) Successful in 5m11s
release / apk x86_64 (push) Failing after 5m8s
release / release apk (push) Has been skipped
release / release (push) Successful in 12s
health plan S7: docs, contract comments, SPEC 019 update
- lx-changelog: health board + observatory + global probe + sub cache entry
- DECISIONS.md: D18 board vs delete-and-overlay, D19 observatory vs sweep, D20 global probe settings
- contract comments: urltest.go CheckOutbounds (fast circuit) + observatory.go loop (background circuit + freshness gate) document the two-circuit split
- SPEC 019: dial-error section updated - slots still not moved, but board verdict demotes dead slot on next pick + retry (§5.B); sticky/replace-in-slot/never-shrink invariants preserved
2026-07-24 18:30:48 +03:00

891 lines
32 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)
}
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,
}
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
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()
}
if s.group.selectedOutboundTCP != nil {
return s.group.selectedOutboundTCP.Tag()
} else if s.group.selectedOutboundUDP != nil {
return s.group.selectedOutboundUDP.Tag()
}
// lx: SPEC 019 — cold start: before the first URL-test, selectedOutbound* is nil 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
selectedOutboundTCP adapter.Outbound
selectedOutboundUDP adapter.Outbound
interruptGroup *interrupt.Group
interruptExternalConnections bool
access sync.Mutex
ticker *time.Ticker
close chan struct{}
started 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
}
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()
g.started = 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).
g.seedPool()
go g.CheckOutbounds(false)
}
func (g *URLTestGroup) Touch() {
if !g.started {
return
}
g.access.Lock()
defer g.access.Unlock()
if g.ticker != nil {
g.lastActive.Store(time.Now())
return
}
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 (g *URLTestGroup) Close() error {
g.access.Lock()
defer g.access.Unlock()
if g.ticker == nil {
return nil
}
g.ticker.Stop()
g.ticker = nil
g.pause.UnregisterCallback(g.pauseCallback)
g.pauseCallback = nil
close(g.close)
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)
}
func (g *URLTestGroup) loopCheck(ticker *time.Ticker, closeChan <-chan struct{}) {
if time.Since(g.lastActive.Load()) > g.interval {
g.lastActive.Store(time.Now())
g.CheckOutbounds(false)
}
for {
select {
case <-closeChan:
return
case <-ticker.C:
}
if 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.CheckOutbounds(false)
}
}
// 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
}
func (g *URLTestGroup) performUpdateCheck() {
var updated bool
if outbound, exists := g.Select(N.NetworkTCP); outbound != nil && (g.selectedOutboundTCP == nil || (exists && outbound != g.selectedOutboundTCP)) {
if g.selectedOutboundTCP != nil {
updated = true
}
g.selectedOutboundTCP = outbound
}
if outbound, exists := g.Select(N.NetworkUDP); outbound != nil && (g.selectedOutboundUDP == nil || (exists && outbound != g.selectedOutboundUDP)) {
if g.selectedOutboundUDP != nil {
updated = true
}
g.selectedOutboundUDP = outbound
}
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
}