feat: HH-773 合集重命名接入 egress geo 缓存、geo 失败降级与每小时补探测
Build and Publish Docker Image / build-and-push (pull_request) Successful in 15m54s

- service: EgressCacheKey 从 handler 下沉复用;MergeCachedEgressGeo 在合集
  重命名前合并缓存 geo 字段(纯缓存读、零网络 I/O);StripEgressGeoFields
  防止 geo 字段泄漏进订阅输出
- filter: geo 检测失败的节点保留 '[别名] 原名'(或纯原名)参与编号/排序,
  不再丢弃
- handler: 每小时 StartEgressRefresher 对启用源补探测缺失/过期的 egress
  缓存(含 TTL 内 error 结果一律跳过);RegisterRoutes 返回 *Deps;
  cmd/server.go 启动定时任务并在 shutdown 时取消
- tests: service 缓存 key 稳定性/合并/管道测试、handler 补探测跳过/TTL/
  过期/取消测试、filter 降级测试;修复 2 个过时测试(mihomo YAML 配置、
  缓存命中路径)
This commit is contained in:
杨豪
2026-08-28 16:40:11 +08:00
parent f3029665ee
commit 5a34e739bb
11 changed files with 543 additions and 49 deletions
+61 -8
View File
@@ -10,6 +10,11 @@ import (
"github.com/peterqiu0516/sub-store/internal/service"
)
// probeNodeEgressFn is the probe seam: a package-level variable so tests can
// stub out the mihomo-based egress probe. Default implementation probes via
// a local mihomo instance.
var probeNodeEgressFn = (*Deps).probeSingleNodeEgress
// probeSourceEgressBackground fetches the source's nodes and probes each
// node's egress info asynchronously. Results are written to the cache so
// that subsequent preview/collection requests can read them without
@@ -24,12 +29,12 @@ func (d *Deps) probeSourceEgressBackground(rec model.SourceRecord) {
settings, _ := d.SettingsRepo.Get()
result, err := service.BuildSubscriptionResult(context.Background(), service.BuildOptions{
Source: &rec,
Sources: []model.SourceRecord{rec},
Target: "json",
Settings: settings,
Source: &rec,
Sources: []model.SourceRecord{rec},
Target: "json",
Settings: settings,
CacheRepo: d.CacheRepo,
ProxyURL: d.Cfg.Fetcher.ProxyURL,
ProxyURL: d.Cfg.Fetcher.ProxyURL,
})
if err != nil {
slog.Warn("background egress probe: failed to build source", "source", rec.ID, "error", err)
@@ -52,15 +57,16 @@ func (d *Deps) probeSourceEgressBackground(rec model.SourceRecord) {
if node == nil {
continue
}
// Skip if already cached
cacheKey := egressCacheKey(node)
// Skip if already cached (fresh entries only — SafeGet treats expired
// entries as misses, and error results are cached with the full TTL too)
cacheKey := service.EgressCacheKey(node)
if _, ok := d.CacheRepo.SafeGet(cacheKey); ok {
probed++
continue
}
// Probe this node
info, err := d.probeSingleNodeEgress(node)
info, err := probeNodeEgressFn(d, node)
if err != nil {
slog.Debug("background egress probe: node failed",
"source", rec.ID, "node", node["name"], "error", err)
@@ -119,3 +125,50 @@ func (d *Deps) probeSingleNodeEgress(node model.ProxyNode) (map[string]any, erro
}
return info, nil
}
// StartEgressRefresher runs a background goroutine that periodically
// re-probes egress info for all enabled sources' nodes whose cache entries
// are missing or expired (HH-773). The first pass runs immediately; fresh
// cache hits — including error results within TTL — are skipped by
// probeSourceEgressBackground. Cancel the context to stop.
func StartEgressRefresher(ctx context.Context, d *Deps, interval time.Duration) {
go func() {
d.refreshAllSourceEgress(ctx)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
d.refreshAllSourceEgress(ctx)
case <-ctx.Done():
return
}
}
}()
}
// refreshAllSourceEgress rebuilds the node list for every enabled source and
// re-probes nodes with missing/expired egress cache entries.
func (d *Deps) refreshAllSourceEgress(ctx context.Context) {
defer func() {
if r := recover(); r != nil {
slog.Warn("egress refresher panicked", "error", r)
}
}()
sources, err := d.SourceRepo.List()
if err != nil {
slog.Warn("egress refresher: failed to list sources", "error", err)
return
}
for _, rec := range sources {
if !rec.Enabled {
continue
}
select {
case <-ctx.Done():
return
default:
}
d.probeSourceEgressBackground(rec)
}
}
+148
View File
@@ -0,0 +1,148 @@
package handler
import (
"context"
"testing"
"time"
"github.com/peterqiu0516/sub-store/internal/model"
"github.com/peterqiu0516/sub-store/internal/proxy"
"github.com/peterqiu0516/sub-store/internal/service"
)
// egressKeyForNode parses a single proxy URI and returns its egress cache key
// (via the service-layer key shared with the pipeline).
func egressKeyForNode(t *testing.T, uri string) string {
t.Helper()
nodes := proxy.ParseProxies(uri)
if len(nodes) != 1 {
t.Fatalf("parse %q: got %d nodes", uri, len(nodes))
}
return service.EgressCacheKey(nodes[0])
}
// stubProbeNodeEgress replaces the mihomo-based probe seam for the duration
// of the test and records every probed node name.
func stubProbeNodeEgress(t *testing.T) *[]string {
t.Helper()
var probed []string
old := probeNodeEgressFn
probeNodeEgressFn = func(d *Deps, node model.ProxyNode) (map[string]any, error) {
probed = append(probed, node["name"].(string))
return map[string]any{"egressIp": "2.2.2.2", "country": "United States", "countryCode": "US", "flag": "🇺🇸"}, nil
}
t.Cleanup(func() { probeNodeEgressFn = old })
return &probed
}
func TestRefreshAllSourceEgress_ProbesMissingSkipsCached(t *testing.T) {
deps := newTestDeps(t)
content := "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#cached-node\nss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.5:81#missing-node"
if _, err := deps.SourceRepo.Upsert(model.SourceRecord{
ID: "s1", Name: "s1", Type: "local", Content: content, Enabled: true,
Filters: []model.FilterRule{}, Meta: map[string]any{},
}); err != nil {
t.Fatalf("upsert source: %v", err)
}
// Disabled sources are never refreshed.
if _, err := deps.SourceRepo.Upsert(model.SourceRecord{
ID: "s2", Name: "s2", Type: "local", Content: "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.6:82#disabled-node", Enabled: false,
Filters: []model.FilterRule{}, Meta: map[string]any{},
}); err != nil {
t.Fatalf("upsert disabled source: %v", err)
}
// Fresh cache entry (geo result) for the first node only.
deps.CacheRepo.SafePut(egressKeyForNode(t, "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#cached-node"),
`{"egressIp":"1.1.1.1","country":"Japan","countryCode":"JP","flag":"🇯🇵"}`, nil, 300)
probed := stubProbeNodeEgress(t)
deps.refreshAllSourceEgress(context.Background())
if len(*probed) != 1 || (*probed)[0] != "missing-node" {
t.Fatalf("probed = %v, want exactly [missing-node] (cached and disabled skipped)", *probed)
}
// The fresh probe result must now be cached.
if _, ok := deps.CacheRepo.SafeGet(egressKeyForNode(t, "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.5:81#missing-node")); !ok {
t.Fatal("probe result should be cached after refresh")
}
}
func TestRefreshAllSourceEgress_SkipsErrorResultsWithinTTL(t *testing.T) {
deps := newTestDeps(t)
if _, err := deps.SourceRepo.Upsert(model.SourceRecord{
ID: "s1", Name: "s1", Type: "local", Content: "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#err-node", Enabled: true,
Filters: []model.FilterRule{}, Meta: map[string]any{},
}); err != nil {
t.Fatalf("upsert source: %v", err)
}
deps.CacheRepo.SafePut(egressKeyForNode(t, "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#err-node"),
`{"egressError":"dial tcp: timeout"}`, nil, 300)
probed := stubProbeNodeEgress(t)
deps.refreshAllSourceEgress(context.Background())
if len(*probed) != 0 {
t.Fatalf("cached error result within TTL must be skipped, probed = %v", *probed)
}
}
func TestRefreshAllSourceEgress_ReprobesExpiredEntries(t *testing.T) {
deps := newTestDeps(t)
if _, err := deps.SourceRepo.Upsert(model.SourceRecord{
ID: "s1", Name: "s1", Type: "local", Content: "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#stale-node", Enabled: true,
Filters: []model.FilterRule{}, Meta: map[string]any{},
}); err != nil {
t.Fatalf("upsert source: %v", err)
}
deps.CacheRepo.SafePut(egressKeyForNode(t, "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#stale-node"),
`{"egressError":"old failure"}`, nil, 300)
// Expire every entry: SafeGet must treat them as misses.
if _, err := deps.DB.Exec("UPDATE source_cache SET cached_at = cached_at - 10000"); err != nil {
t.Fatalf("expire cache entries: %v", err)
}
probed := stubProbeNodeEgress(t)
deps.refreshAllSourceEgress(context.Background())
if len(*probed) != 1 {
t.Fatalf("expired entry must be re-probed, probed = %v", *probed)
}
if _, ok := deps.CacheRepo.SafeGet(egressKeyForNode(t, "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#stale-node")); !ok {
t.Fatal("fresh result should be re-cached after re-probe")
}
}
// StartEgressRefresher must run the first refresh immediately and stop when
// the context is cancelled.
func TestStartEgressRefresher_ImmediateRunAndCancel(t *testing.T) {
deps := newTestDeps(t)
if _, err := deps.SourceRepo.Upsert(model.SourceRecord{
ID: "s1", Name: "s1", Type: "local", Content: "ss://YWVzLTI1Ni1nY206cGFzcw==@1.2.3.4:80#n1", Enabled: true,
Filters: []model.FilterRule{}, Meta: map[string]any{},
}); err != nil {
t.Fatalf("upsert source: %v", err)
}
done := make(chan string, 8)
old := probeNodeEgressFn
probeNodeEgressFn = func(d *Deps, node model.ProxyNode) (map[string]any, error) {
done <- node["name"].(string)
return map[string]any{"country": "United States"}, nil
}
t.Cleanup(func() { probeNodeEgressFn = old })
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
StartEgressRefresher(ctx, deps, time.Hour)
select {
case name := <-done:
if name != "n1" {
t.Fatalf("probed %q, want n1", name)
}
case <-time.After(5 * time.Second):
t.Fatal("refresher did not run its first pass immediately")
}
cancel()
}
+3 -21
View File
@@ -3,7 +3,6 @@ package handler
import (
"bytes"
"context"
"crypto/sha256"
"encoding/json"
"fmt"
"io"
@@ -20,6 +19,7 @@ import (
"gopkg.in/yaml.v3"
"github.com/peterqiu0516/sub-store/internal/model"
"github.com/peterqiu0516/sub-store/internal/service"
"github.com/peterqiu0516/sub-store/internal/util"
)
@@ -35,7 +35,7 @@ func (d *Deps) HandleEgressInfo(c fiber.Ctx) error {
node["name"] = "PROXY"
}
cacheKey := egressCacheKey(node)
cacheKey := service.EgressCacheKey(node)
if entry, ok := d.CacheRepo.SafeGet(cacheKey); ok {
var cached map[string]any
if json.Unmarshal([]byte(entry.Content), &cached) == nil {
@@ -75,27 +75,9 @@ func (d *Deps) HandleEgressInfo(c fiber.Ctx) error {
return success(c, info)
}
func egressCacheKey(node model.ProxyNode) string {
clean := model.ProxyNode{}
skip := map[string]bool{
"id": true, "name": true, "_sourceAlias": true, "_previewId": true,
"latencyMs": true, "latencyError": true,
"egressIp": true, "egressCountry": true, "egressRegion": true, "egressError": true,
"country": true, "countryCode": true, "region": true, "city": true, "isp": true, "flag": true, "cached": true,
}
for k, v := range node {
if !skip[k] {
clean[k] = v
}
}
data, _ := json.Marshal(clean)
sum := sha256.Sum256(data)
return fmt.Sprintf("egress:%x", sum)
}
func (d *Deps) addCachedEgressInfo(nodes []model.ProxyNode) []model.ProxyNode {
for _, node := range nodes {
entry, ok := d.CacheRepo.SafeGet(egressCacheKey(node))
entry, ok := d.CacheRepo.SafeGet(service.EgressCacheKey(node))
if !ok {
continue
}
+31 -12
View File
@@ -14,11 +14,13 @@ import (
"github.com/gofiber/fiber/v3"
"github.com/jmoiron/sqlx"
"gopkg.in/yaml.v3"
_ "modernc.org/sqlite"
"github.com/peterqiu0516/sub-store/internal/config"
"github.com/peterqiu0516/sub-store/internal/database"
"github.com/peterqiu0516/sub-store/internal/model"
"github.com/peterqiu0516/sub-store/internal/service"
"github.com/peterqiu0516/sub-store/internal/template"
)
@@ -2127,25 +2129,42 @@ func TestBuildEgressProbeConfig(t *testing.T) {
if err != nil {
t.Fatalf("buildEgressProbeConfig returned error: %v", err)
}
// The probe config is a minimal mihomo (Clash Meta) YAML.
var doc map[string]any
if err := json.Unmarshal(data, &doc); err != nil {
t.Fatalf("invalid config JSON: %v", err)
if err := yaml.Unmarshal(data, &doc); err != nil {
t.Fatalf("invalid config YAML: %v", err)
}
route := doc["route"].(map[string]any)
if route["final"] != "PROXY" {
t.Fatalf("route.final = %v, want PROXY", route["final"])
if port, _ := doc["mixed-port"].(int); port != 19090 {
t.Fatalf("mixed-port = %v, want 19090", doc["mixed-port"])
}
inbound := doc["inbounds"].([]any)[0].(map[string]any)
if inbound["listen"] != "127.0.0.1" || inbound["listen_port"].(float64) != 19090 {
t.Fatalf("unexpected inbound: %v", inbound)
proxies, _ := doc["proxies"].([]any)
if len(proxies) != 1 {
t.Fatalf("proxies = %v, want 1 entry", doc["proxies"])
}
proxyMap, _ := proxies[0].(map[string]any)
if proxyMap["name"] != "PROXY" || proxyMap["server"] != "127.0.0.1" || proxyMap["cipher"] != "aes-256-gcm" {
t.Fatalf("unexpected probe proxy: %v", proxyMap)
}
groups, _ := doc["proxy-groups"].([]any)
group, _ := groups[0].(map[string]any)
if group["type"] != "select" {
t.Fatalf("proxy-group type = %v, want select", group["type"])
}
}
func TestHandleEgressInfoUnsupported(t *testing.T) {
func TestHandleEgressInfoCacheHit(t *testing.T) {
deps := newTestDeps(t)
app := newApp(deps)
code, _ := doRequest(t, app, "POST", "/api/utils/egress-info", `{"type":"unknown","server":"1.2.3.4","port":443}`, nil)
assertStatus(t, "EgressInfo unsupported", code, 400)
// Any node type is probed now; a cached result must short-circuit before
// any probing so the response is deterministic without network access.
deps.CacheRepo.SafePut(service.EgressCacheKey(model.ProxyNode{"type": "ss", "server": "1.2.3.4", "port": 8388}),
`{"egressIp":"1.1.1.1","country":"Japan","latencyMs":12}`, nil, 300)
code, out := doRequest(t, app, "POST", "/api/utils/egress-info", `{"type":"ss","server":"1.2.3.4","port":8388}`, nil)
assertStatus(t, "EgressInfo cache hit", code, 200)
data, _ := out["data"].(map[string]any)
if data["egressIp"] != "1.1.1.1" || data["cached"] != true {
t.Fatalf("unexpected cached egress response: %v", data)
}
}
func TestProbeServerPortLatency(t *testing.T) {
@@ -2180,7 +2199,7 @@ func TestAddCachedEgressInfo(t *testing.T) {
"server": "127.0.0.1",
"port": 8388,
}
deps.CacheRepo.SafePut(egressCacheKey(node), `{"egressIp":"1.1.1.1","country":"Japan","region":"Tokyo","latencyMs":12}`, nil, 300)
deps.CacheRepo.SafePut(service.EgressCacheKey(node), `{"egressIp":"1.1.1.1","country":"Japan","region":"Tokyo","latencyMs":12}`, nil, 300)
nodes := deps.addCachedEgressInfo([]model.ProxyNode{node})
if nodes[0]["egressIp"] != "1.1.1.1" || nodes[0]["latencyMs"].(float64) != 12 {
t.Fatalf("cached egress not merged: %v", nodes[0])
+4 -2
View File
@@ -34,8 +34,9 @@ func NewDeps(cfg *config.Config, db *sqlx.DB) *Deps {
}
}
// RegisterRoutes registers all API and download routes.
func RegisterRoutes(app *fiber.App, cfg *config.Config, db *sqlx.DB) {
// RegisterRoutes wires all routes and returns the Deps used, so callers
// (e.g. server startup) can launch background jobs on the same dependencies.
func RegisterRoutes(app *fiber.App, cfg *config.Config, db *sqlx.DB) *Deps {
deps := NewDeps(cfg, db)
// Admin API group — requires admin token
@@ -103,6 +104,7 @@ func RegisterRoutes(app *fiber.App, cfg *config.Config, db *sqlx.DB) {
// Public download routes — no admin token required, uses download token
app.Get("/sources/:name/:token", deps.HandleDownloadSource)
app.Get("/collections/:name/:token", deps.HandleDownloadCollection)
return deps
}
// success sends a success JSON response.