fix(overview): stop the "everything is down" flash when opening Advanced mode
Root cause was performance, not the toggle. NodesJSON TCP-dials every node to
fill liveness/latency, so on a full subscription (258+ nodes) one call took ~7s.
The Overview polls the node list every 5s, so the dials stacked and starved the
box; the Advanced switch does a full page reload, and under that load a slow/
empty `status` RPC would paint every tile red until the next good poll.
Fix, three layers:
- xrayctl: cache the probe sweep on disk (/var/run/shater/nodes_probe.json,
45s TTL) and single-flight the refresh with a NON-BLOCKING lock, so the hot
path is a file read (7s -> 0.03s) and two callers never dial at once — a
late caller serves the cache instead of piling on. `nodes probe` forces a
fresh sweep; the dashboard's bare `nodes` serves the cache.
- ubus: the `nodes` method takes an optional `probe` arg (default false).
- overview.js: a failed/timed-out status RPC resolves to {} — keep the last
good view instead of flashing "down", and let "updated Xs ago" show staleness.
As a bonus the Nodes page latency column now populates from the cache without a
manual "Test all". Verified live: warm `nodes` 0.03s, single-flight confirmed
(one sweep, concurrent call returns immediately), Advanced toggle loads clean
and green with 0 console errors; `xray -test` OK on a config built from 337
real-subscription nodes. Bumps xrayctl r19, luci-app-shater r9.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LLthkP2S8WAfxu7fcYbPfE
This commit is contained in:
@@ -11,7 +11,7 @@ LUCI_PKGARCH:=all
|
||||
# Explicit version so the data-.ipk packer (ci/pack-luci.sh) and any SDK build
|
||||
# agree; bump PKG_RELEASE to ship a UI upgrade via `opkg upgrade`.
|
||||
PKG_VERSION:=0.1.0
|
||||
PKG_RELEASE:=8
|
||||
PKG_RELEASE:=9
|
||||
|
||||
PKG_LICENSE:=GPL-2.0
|
||||
PKG_MAINTAINER:=shater <maqrota@icloud.com>
|
||||
|
||||
@@ -679,6 +679,12 @@ return view.extend({
|
||||
L.resolveDefault(callNodes(), {}),
|
||||
L.resolveDefault(callSubInfo(), {})
|
||||
]).then(function(res) {
|
||||
// A failed/timed-out status RPC resolves to {} (StatusJSON always
|
||||
// returns a populated object otherwise). Don't flash "everything
|
||||
// down" on a transient hiccup — keep the last good view and let the
|
||||
// "updated Xs ago" tick reveal the staleness.
|
||||
if (!res[0] || typeof res[0] != 'object' || !Object.keys(res[0]).length)
|
||||
return;
|
||||
var m = computeModel(res[0], res[1], res[2]);
|
||||
dom.content(pathBox, renderSignalPath(m));
|
||||
dom.content(kpiBox, renderKPIs(m));
|
||||
|
||||
@@ -45,8 +45,10 @@ function shq(v) {
|
||||
}
|
||||
|
||||
// Normalize xrayctl `nodes` output (a JSON array) into { nodes: [...] }.
|
||||
function nodes_result() {
|
||||
let d = run_json(XRAYCTL + ' nodes');
|
||||
// probe=true forces a fresh liveness sweep; the default serves xrayctl's cached
|
||||
// sweep, which keeps the dashboard's frequent polling cheap.
|
||||
function nodes_result(probe) {
|
||||
let d = run_json(XRAYCTL + ' nodes' + (probe ? ' probe' : ''));
|
||||
if (type(d) == 'array')
|
||||
return { nodes: d };
|
||||
return d;
|
||||
@@ -68,8 +70,9 @@ return {
|
||||
},
|
||||
|
||||
nodes: {
|
||||
call: function() {
|
||||
return nodes_result();
|
||||
args: { probe: false },
|
||||
call: function(req) {
|
||||
return nodes_result(req.args?.probe);
|
||||
}
|
||||
},
|
||||
|
||||
|
||||
+1
-1
@@ -10,7 +10,7 @@ include $(TOPDIR)/rules.mk
|
||||
|
||||
PKG_NAME:=xrayctl
|
||||
PKG_VERSION:=0.1.0
|
||||
PKG_RELEASE:=18
|
||||
PKG_RELEASE:=19
|
||||
|
||||
PKG_MAINTAINER:=shater
|
||||
PKG_LICENSE:=MIT
|
||||
|
||||
+22
-1
@@ -17,7 +17,28 @@ import (
|
||||
|
||||
const lockPath = "/var/lock/xrayctl.lock"
|
||||
|
||||
func init() { lockExclusive = flockExclusive }
|
||||
func init() { lockExclusive = flockExclusive; lockTry = flockTryPath }
|
||||
|
||||
// flockTryPath is the non-blocking counterpart used to single-flight the node
|
||||
// probe sweep. ok=false (EWOULDBLOCK) means another process already holds it —
|
||||
// that is expected contention, not an error, so err stays nil.
|
||||
func flockTryPath(path string) (release func(), ok bool, err error) {
|
||||
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
f, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0o644)
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
if err := syscall.Flock(int(f.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); err != nil {
|
||||
_ = f.Close()
|
||||
return nil, false, nil
|
||||
}
|
||||
return func() {
|
||||
_ = syscall.Flock(int(f.Fd()), syscall.LOCK_UN)
|
||||
_ = f.Close()
|
||||
}, true, nil
|
||||
}
|
||||
|
||||
func flockExclusive() (release func(), err error) {
|
||||
if err := os.MkdirAll(filepath.Dir(lockPath), 0o755); err != nil {
|
||||
|
||||
+9
-1
@@ -24,6 +24,12 @@ func main() {
|
||||
// flock_unix.go's init installs the real implementation on the target OS.
|
||||
var lockExclusive = func() (release func(), err error) { return func() {}, nil }
|
||||
|
||||
// lockTry attempts a NON-BLOCKING cross-process lock on `path`: ok=false means
|
||||
// another process already holds it, so the caller should back off rather than
|
||||
// wait (used to single-flight the expensive node probe sweep). The no-op default
|
||||
// (non-unix dev/test hosts) always "succeeds".
|
||||
var lockTry = func(path string) (release func(), ok bool, err error) { return func() {}, true, nil }
|
||||
|
||||
// mutatingCmd reports whether a command mutates shared state (run.json, nft
|
||||
// table, policy routing, UCI) and therefore must hold the exclusive lock.
|
||||
func mutatingCmd(cmd string, rest []string) bool {
|
||||
@@ -97,7 +103,9 @@ func run(args []string) int {
|
||||
case "conns":
|
||||
return emit(ConnsJSON())
|
||||
case "nodes":
|
||||
return emit(NodesJSON())
|
||||
// `nodes probe` forces a fresh liveness sweep; bare `nodes` serves the
|
||||
// cached sweep (fast — the dashboard polls this every few seconds).
|
||||
return emit(NodesJSON(len(rest) > 0 && rest[0] == "probe"))
|
||||
case "stats":
|
||||
return emit(StatsJSON())
|
||||
case "explain":
|
||||
|
||||
+104
-23
@@ -43,10 +43,84 @@ func StatusJSON() ([]byte, error) {
|
||||
return json.MarshalIndent(st, "", " ")
|
||||
}
|
||||
|
||||
// Probe-liveness cache. NodesJSON's per-endpoint TCP dial is the expensive part:
|
||||
// a full subscription is easily hundreds of nodes and each dial blocks up to the
|
||||
// stProbe timeout, so one sweep takes several seconds. The dashboard polls the
|
||||
// node list every few seconds, so dialing on every call would stack multi-second
|
||||
// invocations and starve the router (this was the "everything looks down" flash
|
||||
// after a page reload). Instead the last sweep is cached on disk and refreshed at
|
||||
// most once per nodesProbeTTL, single-flighted with a non-blocking lock so a slow
|
||||
// sweep is never run by two callers at once — late callers just serve the cache.
|
||||
var (
|
||||
nodesProbeCache = "/var/run/shater/nodes_probe.json"
|
||||
nodesProbeLock = "/var/run/shater/nodes_probe.lock"
|
||||
)
|
||||
|
||||
const nodesProbeTTL = 45 * time.Second
|
||||
|
||||
type probeRec struct {
|
||||
Alive bool `json:"alive"`
|
||||
MS int `json:"ms"`
|
||||
}
|
||||
type probeCacheFile struct {
|
||||
TS int64 `json:"ts"`
|
||||
Nodes map[string]probeRec `json:"nodes"` // keyed by node name
|
||||
}
|
||||
|
||||
func loadProbeCache() probeCacheFile {
|
||||
var c probeCacheFile
|
||||
if b, err := os.ReadFile(nodesProbeCache); err == nil {
|
||||
_ = json.Unmarshal(b, &c)
|
||||
}
|
||||
if c.Nodes == nil {
|
||||
c.Nodes = map[string]probeRec{}
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
func saveProbeCache(c probeCacheFile) {
|
||||
if os.MkdirAll("/var/run/shater", 0o755) != nil {
|
||||
return
|
||||
}
|
||||
if b, err := json.Marshal(c); err == nil {
|
||||
tmp := nodesProbeCache + ".tmp"
|
||||
if os.WriteFile(tmp, b, 0o644) == nil {
|
||||
_ = os.Rename(tmp, nodesProbeCache)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// sweepProbes dials every node's endpoint (bounded fan-out) and returns a fresh
|
||||
// liveness cache. addrs[i]/ports[i] correspond to names[i].
|
||||
func sweepProbes(names, addrs []string, ports []int) probeCacheFile {
|
||||
c := probeCacheFile{Nodes: make(map[string]probeRec, len(names))}
|
||||
var mu sync.Mutex
|
||||
sem := make(chan struct{}, 32)
|
||||
var wg sync.WaitGroup
|
||||
for i := range names {
|
||||
wg.Add(1)
|
||||
sem <- struct{}{}
|
||||
go func(name, a string, p int) {
|
||||
defer wg.Done()
|
||||
defer func() { <-sem }()
|
||||
alive, ms, _ := stProbe(a, p)
|
||||
mu.Lock()
|
||||
c.Nodes[name] = probeRec{Alive: alive, MS: ms}
|
||||
mu.Unlock()
|
||||
}(names[i], addrs[i], ports[i])
|
||||
}
|
||||
wg.Wait()
|
||||
c.TS = time.Now().Unix()
|
||||
return c
|
||||
}
|
||||
|
||||
// NodesJSON lists all nodes (manual + subscription cache) with liveness fields.
|
||||
// Liveness is merged from xray's API (observatory-selected member + traffic
|
||||
// counters) and, for latency and as a fallback, a real endpoint probe.
|
||||
func NodesJSON() ([]byte, error) {
|
||||
// counters, always fresh) and, for latency and as a fallback, a cached endpoint
|
||||
// probe (see the probe cache above). force=true runs a fresh sweep regardless of
|
||||
// cache age (the `nodes probe` subcommand); the dashboard's plain `nodes` call
|
||||
// serves the cache and stays fast.
|
||||
func NodesJSON(force bool) ([]byte, error) {
|
||||
m, err := ReadUCI()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -71,9 +145,11 @@ func NodesJSON() ([]byte, error) {
|
||||
URI string `json:"uri"` // share-link (for QR / link export)
|
||||
}
|
||||
|
||||
// Cheap pass: identity + observatory/stats-derived liveness (no dialing).
|
||||
out := make([]nodeOut, len(m.Nodes))
|
||||
sem := make(chan struct{}, 32)
|
||||
var wg sync.WaitGroup
|
||||
addrs := make([]string, len(m.Nodes))
|
||||
ports := make([]int, len(m.Nodes))
|
||||
names := make([]string, len(m.Nodes))
|
||||
for i, n := range m.Nodes {
|
||||
proto, fp := "", n.Fingerprint
|
||||
var addr string
|
||||
@@ -85,6 +161,7 @@ func NodesJSON() ([]byte, error) {
|
||||
}
|
||||
addr, port, _, _, _, _, _ = extractIdentity(ob)
|
||||
}
|
||||
names[i], addrs[i], ports[i] = n.Name, addr, port
|
||||
tag := tags[n.Name]
|
||||
traffic := obs.Stats[tag]
|
||||
selected := tag != "" && obs.Selected[tag]
|
||||
@@ -102,26 +179,30 @@ func NodesJSON() ([]byte, error) {
|
||||
rec.Alive, rec.Source = true, "stats"
|
||||
}
|
||||
out[i] = rec
|
||||
|
||||
// Endpoint probe (latency, and alive fallback when the API is silent).
|
||||
// Bounded like NodeTest: an unbounded fan-out over a huge subscription
|
||||
// cache would open hundreds of concurrent dials on a small router.
|
||||
wg.Add(1)
|
||||
sem <- struct{}{}
|
||||
go func(idx int, a string, p int) {
|
||||
defer wg.Done()
|
||||
defer func() { <-sem }()
|
||||
alive, ms, _ := stProbe(a, p)
|
||||
out[idx].LatencyMS = ms
|
||||
if out[idx].Source == "" {
|
||||
if alive {
|
||||
out[idx].Alive = true
|
||||
}
|
||||
out[idx].Source = "probe"
|
||||
}
|
||||
}(i, addr, port)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
// Probe layer: refresh the cached sweep at most once per TTL (or when forced),
|
||||
// single-flighted so a slow sweep never stacks; otherwise reuse the cache.
|
||||
cache := loadProbeCache()
|
||||
if force || cache.TS == 0 || time.Since(time.Unix(cache.TS, 0)) > nodesProbeTTL {
|
||||
if release, ok, _ := lockTry(nodesProbeLock); ok {
|
||||
cache = sweepProbes(names, addrs, ports)
|
||||
saveProbeCache(cache)
|
||||
release()
|
||||
}
|
||||
// !ok: another sweep is already in flight — fall through with the old cache.
|
||||
}
|
||||
for i := range out {
|
||||
if pr, ok := cache.Nodes[out[i].Name]; ok {
|
||||
out[i].LatencyMS = pr.MS
|
||||
if out[i].Source == "" {
|
||||
if pr.Alive {
|
||||
out[i].Alive = true
|
||||
}
|
||||
out[i].Source = "probe"
|
||||
}
|
||||
}
|
||||
}
|
||||
return json.MarshalIndent(out, "", " ")
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user