// manager.go:ProxyManager 装配与出口决策核心(design §5)。 // // 职责:订阅解析 → provider 落盘 + mihomo 热载 → 探活循环(30s/5min 双频)→ // /api/exit 决策(deny 规则 → sticky → P2C 主备区域组)。 // fail-closed:订阅拉取失败且无缓存 provider 时,所有 exit 请求 unhealthy。 package proxymanager import ( "context" "fmt" "log" "math/rand" "os" "sync" "time" "onesvm.com/onesvm/browser-server/internal/config" "onesvm.com/onesvm/browser-server/internal/policy" "onesvm.com/onesvm/browser-server/internal/store" ) // mixedProxyURL 统一 mixed 出口(D3:mihomo mixed :17890 仅 overlay)。 const mixedProxyURL = "http://mihomo:17890" // Manager ProxyManager 控制面。 type Manager struct { mu sync.RWMutex health *HealthEngine sel *Selector subs *subscriptionCache ctrl *controller rules *policy.DomainTrie db *store.DB cfgDir string logger *log.Logger activeNode string // 当前活跃出口(provenance/healthz) lastSwitch time.Time // 最近切换时刻 lastRefresh time.Time // 最近一次订阅成功刷新 healthy bool // 出口是否可用(订阅+provider 就绪且池存活) unhealthy string // unhealthy 原因 subStale bool // 订阅 stale-on-error 状态 switches int } // NewManager 构造。db 可 nil(单测无 SQLite 场景,路由表仅内存默认)。 func NewManager(db *store.DB, cfgDir string, subURLs []string, ctrlURL, ctrlSecret string, logger *log.Logger) *Manager { if logger == nil { logger = log.New(os.Stdout, "[browser-server/proxymanager] ", log.LstdFlags) } ctrl := newController(ctrlURL, ctrlSecret, 10*time.Second) return &Manager{ health: NewHealthEngine(ctrl, logger), // W6:生产探活= mihomo delay API;禁 nil sel: NewSelector(), subs: newSubscriptionCache(subURLs), ctrl: ctrl, rules: policy.NewDomainTrie(), db: db, cfgDir: cfgDir, logger: logger, } } // SetProber 注入探针实现(生产= controller;测试= stub)。 func (m *Manager) SetProber(p Prober) { m.health.prober = p } // Refresh 拉订阅 → 解析 → 写 provider → mihomo 热载 → 重建探活池。 // 订阅拉取失败且无缓存:保持现有池不变并置 unhealthy(fail-closed); // 有缓存(stale-on-error):续用旧配置并告警。 func (m *Manager) Refresh(ctx context.Context) error { text, stale, err := m.subs.Get(ctx) if err != nil { m.mu.Lock() m.healthy = false m.unhealthy = "订阅不可用(fail-closed)" m.subStale = false m.mu.Unlock() return fmt.Errorf("manager: %w", err) } parsed, err := ParseSubscription(string(text)) if err != nil { m.mu.Lock() m.healthy = false m.unhealthy = "订阅解析失败: " + err.Error() m.subStale = stale m.mu.Unlock() return fmt.Errorf("manager: %w", err) } cfg, err := BuildProviderConfig(string(text)) if err != nil { m.mu.Lock() m.healthy = false m.unhealthy = "provider 生成失败: " + err.Error() m.subStale = stale m.mu.Unlock() return fmt.Errorf("manager: %w", err) } if err := WriteProvider(m.cfgDir, cfg); err != nil { m.mu.Lock() m.healthy = false m.unhealthy = "provider 落盘失败: " + err.Error() m.mu.Unlock() return fmt.Errorf("manager: %w", err) } // mihomo 热载:不可达仅告警不 panic(mihomo 可能由 stack 另行编排启动)。 if err := m.ctrl.Reload(ctx, MihomoReloadPath()); err != nil { m.logger.Printf("[manager] 告警: mihomo 热载失败(不阻塞): %v", err) } // 探活池重建(保留既有状态机进度)+ 路由表热载。 m.health.ReplacePool(parsed.ProxiesMeta) m.loadRules() m.mu.Lock() m.healthy = true m.unhealthy = "" m.subStale = stale m.lastRefresh = config.Now() m.mu.Unlock() s := parsed.Summary() m.logger.Printf("[manager] 订阅刷新 format=%s raw=%d real=%d vless=%d hy2=%d regions=%v stale=%v", parsed.Format, parsed.RawCount, len(parsed.ProxiesMeta), len(parsed.Pool[poolVless]), len(parsed.Pool[poolUDP]), SortedRegions(parsed.RegionCount(poolVless)), stale) _ = s // 摘要仅供调试挂点;计数已入日志(脱敏:无 URL 无凭据) return nil } // loadRules 从 store rules 表热载域名路由(db 为 nil 时仅保留默认)。 func (m *Manager) loadRules() { if m.db == nil { return } if _, err := m.rules.LoadFromStore(m.db); err != nil { m.logger.Printf("[manager] 告警: 路由表加载失败(沿用旧表): %v", err) } } // ReloadRules /api/rules/reload 入口。 func (m *Manager) ReloadRules() error { m.loadRules() return nil } // exitRegionFor 按主备区域组策略给 exit 的期望区域: // 优先美国组(圣何塞实测最优),备日本组(防美国入口集体抖动); // 同区域组内 P2C 决定具体节点(节点名不硬编码)。 func exitRegionFor(candidates map[string]int) string { for _, want := range RegionPrefixPool { if candidates[want] > 0 { return want } } // 两组皆无节点:回退任一候选区域。 for _, r := range SortedRegions(candidates) { if candidates[r] > 0 { return r } } return "" } // GetExit /api/exit 决策:deny 规则 → sticky → P2C 主备区域组。 // domain/session 至少其一非空(handler 层校验)。 func (m *Manager) GetExit(domain, session string) ExitDecision { host := domain m.mu.RLock() healthy, unhealthyReason := m.healthy, m.unhealthy m.mu.RUnlock() if !healthy { return ExitDecision{Proxy: mixedProxyURL, Blocked: true, Reason: "unhealthy:" + unhealthyReason} } // deny 域:blocked=true,gateway/scheduler 侧据此拒绝。 if action, ok := m.rules.Lookup(host); ok && action == policy.ActionDeny { return ExitDecision{Proxy: mixedProxyURL, Blocked: true, Reason: "deny_rule"} } // sticky 命中:直接回钉死出口(TTL 内同一 session/domain 钉死)。 key := StickyKey(session, host) if node, ok := m.sel.sticky.Get(key); ok { return ExitDecision{Proxy: mixedProxyURL, Node: node, Region: m.health.RegionOf(node), Sticky: true} } // 区域组内选候选:主区域组优先,空则全池兜底。 candAll := m.health.Candidates() if len(candAll) == 0 { return ExitDecision{Proxy: mixedProxyURL, Blocked: true, Reason: "pool_empty"} } byRegion := map[string][]string{} for _, n := range candAll { byRegion[m.health.RegionOf(n)] = append(byRegion[m.health.RegionOf(n)], n) } region := exitRegionFor(countRegions(byRegion)) pool := byRegion[region] if len(pool) == 0 { pool = candAll } node := Pick(pool, m.health.NodeDelay, func(n int) int { return rand.Intn(n) }) if node == "" { return ExitDecision{Proxy: mixedProxyURL, Blocked: true, Reason: "pool_empty"} } m.sel.sticky.Set(key, node) m.mu.Lock() if m.activeNode != node { m.activeNode = node m.lastSwitch = config.Now() m.switches++ } m.mu.Unlock() // mihomo 层把 VLESS 组钉到该节点(尽力而为;失败不阻塞决策—— // mixed 出口仍可用,仅可能暂与 provenance 显示不一致)。 go func() { ctx2, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if err := m.ctrl.SelectInGroup(ctx2, selectorVless, node); err != nil { m.logger.Printf("[manager] 告警: 选择组切换失败 node=%s err=%v", node, err) } }() return ExitDecision{Proxy: mixedProxyURL, Node: node, Region: m.health.RegionOf(node)} } // countRegions 区域→节点数。 func countRegions(byRegion map[string][]string) map[string]int { out := map[string]int{} for r, nodes := range byRegion { out[r] = len(nodes) } return out } // Healthz /healthz 数据。 func (m *Manager) Healthz() map[string]any { m.mu.RLock() defer m.mu.RUnlock() return map[string]any{ "ok": m.healthy && m.health.PoolAlive(), "pool_alive": m.health.PoolAlive(), "pool_total": m.health.TotalCount(), "active_exit": m.activeNode, "last_switch": contractTimeStr(m.lastSwitch), "last_refresh": contractTimeStr(m.lastRefresh), "sub_stale": m.subStale, "switches": m.switches, "unhealthy": m.unhealthy, } } // Proxies /api/proxies 数据(探活状态,脱敏无凭据)。 func (m *Manager) Proxies() []NodeSnapshot { return m.health.Snapshot() } // Run 后台循环:订阅刷新(TTL 内跳过)+ 活跃 30s + 全池 5min 探活。 // 阻塞至 ctx 取消。 func (m *Manager) Run(ctx context.Context) { // 启动即首次刷新(fail-closed:失败不 panic,healthz 报 unhealthy)。 if err := m.Refresh(ctx); err != nil { m.logger.Printf("[manager] 首次订阅刷新失败: %v", err) } subTick := time.NewTicker(subTTLCache) defer subTick.Stop() activeTick := time.NewTicker(activeInterval) defer activeTick.Stop() fullTick := time.NewTicker(fullPoolInterval) defer fullTick.Stop() for { select { case <-ctx.Done(): return case <-subTick.C: if err := m.Refresh(ctx); err != nil { m.logger.Printf("[manager] 订阅刷新失败: %v", err) } case <-activeTick.C: m.mu.RLock() active := m.activeNode m.mu.RUnlock() if active == "" { continue } ok, failed := m.health.ProbeActive(ctx, []string{active}) if !ok { m.logger.Printf("[manager] 活跃出口探活失败 node=%s failed=%v", active, failed) // 失败立即重选(sticky 不清——同键下次到期重选)。 m.sel.sticky.Reset() } case <-fullTick.C: alive, removed := m.health.ProbeFullPool(ctx) m.logger.Printf("[manager] 全池探活完成 alive=%d removed=%d total=%d", alive, removed, m.health.TotalCount()) } } } // contractTimeStr 东八区 RFC3339(零值给空串)。 func contractTimeStr(t time.Time) string { if t.IsZero() { return "" } return t.In(config.TZ).Format(time.RFC3339) }