package core import ( "context" "errors" "fmt" "io" "net" "net/http" "net/url" "os" "strings" "sync" "time" "vpnclient/internal/config" "vpnclient/internal/corebin" "vpnclient/internal/linknorm" "vpnclient/internal/protocols/awg" "vpnclient/internal/protocols/hysteria2" "vpnclient/internal/protocols/naive" "vpnclient/internal/protocols/xray" "vpnclient/internal/subscription" "vpnclient/internal/sysproxy" ) // Manager orchestrates protocol engines and OS system proxy. type Manager struct { mu sync.Mutex cfgPath string cfg *config.Config engine Engine sys sysproxy.Controller profile *config.Profile stderr io.Writer binDir string // lastStop tracks in-flight engine.Stop so Connect waits for ports to free. lastStop sync.WaitGroup lastError string unexpectedDrop bool watchGeneration uint64 } func NewManager(cfgPath string, cfg *config.Config, stderr io.Writer) (*Manager, error) { if stderr == nil { stderr = os.Stderr } binDir, err := config.ResolveBinDir(cfgPath, cfg.BinDir) if err != nil { return nil, err } m := &Manager{ cfgPath: cfgPath, cfg: cfg, sys: sysproxy.New(), stderr: stderr, binDir: binDir, } m.RecoverSystemProxy() return m, nil } // RecoverSystemProxy clears a system proxy left on after an unclean exit. func (m *Manager) RecoverSystemProxy() { if m == nil || m.sys == nil { return } if !sysproxy.HasSentinel(m.cfgPath) { return } if err := m.sys.ForceDisable(); err != nil { fmt.Fprintf(m.stderr, "recover system proxy: %v\n", err) } sysproxy.ClearSentinel(m.cfgPath) } func (m *Manager) waitEngineStopped() { m.lastStop.Wait() } func (m *Manager) Connect(ctx context.Context, profileName string) error { // Wait outside the lock so a slow Disconnect can finish freeing :1080/:1081. m.waitEngineStopped() m.mu.Lock() if m.engine != nil && m.engine.Running() { m.mu.Unlock() return fmt.Errorf("already connected; disconnect first") } if m.engine != nil { m.engine = nil m.profile = nil } if profileName != "" { m.cfg.Active = profileName } profile, err := m.cfg.ActiveProfile() if err != nil { m.mu.Unlock() return err } profileCopy := *profile binDir := m.binDir wantSysProxy := m.cfg.SystemProxy cfgPath := m.cfgPath sys := m.sys stderr := m.stderr m.mu.Unlock() engine, err := m.newEngine(profileCopy.Protocol) if err != nil { return err } if err := engine.Start(ctx, profileCopy, binDir); err != nil { return err } m.mu.Lock() if m.engine != nil && m.engine.Running() { _ = engine.Stop() m.mu.Unlock() return fmt.Errorf("already connected; disconnect first") } if wantSysProxy { httpProxy, ok := engine.LocalHTTPProxy() if !ok { m.mu.Unlock() _ = engine.Stop() return fmt.Errorf("system_proxy requires an http:// listen address in the profile") } if err := sys.Enable(httpProxy); err != nil { if errors.Is(err, sysproxy.ErrUnsupported) { fmt.Fprintf(stderr, "system proxy unsupported on this OS; local SOCKS/HTTP still up\n") } else { m.mu.Unlock() _ = engine.Stop() return fmt.Errorf("enable system proxy: %w", err) } } else { sysproxy.WriteSentinel(cfgPath, httpProxy) } } m.engine = engine m.profile = &profileCopy m.lastError = "" m.unexpectedDrop = false m.watchGeneration++ gen := m.watchGeneration m.mu.Unlock() go m.watchEngine(gen, engine) return nil } // LastError returns the last connection / core failure message. func (m *Manager) LastError() string { m.mu.Lock() defer m.mu.Unlock() return m.lastError } // SetLastError stores a UI-facing error string. func (m *Manager) SetLastError(msg string) { m.mu.Lock() m.lastError = strings.TrimSpace(msg) m.mu.Unlock() } // TakeUnexpectedDrop reports and clears a crash-driven disconnect. func (m *Manager) TakeUnexpectedDrop() bool { m.mu.Lock() defer m.mu.Unlock() if !m.unexpectedDrop { return false } m.unexpectedDrop = false return true } func (m *Manager) watchEngine(gen uint64, eng Engine) { if eng == nil { return } for { time.Sleep(800 * time.Millisecond) m.mu.Lock() if m.watchGeneration != gen || m.engine != eng { m.mu.Unlock() return } if eng.Running() { m.mu.Unlock() continue } // Core died without Disconnect(). m.engine = nil m.profile = nil m.lastError = "процесс ядра завершился неожиданно" m.unexpectedDrop = true sys := m.sys cfgPath := m.cfgPath m.mu.Unlock() if sys != nil && (sys.Enabled() || sysproxy.HasSentinel(cfgPath)) { _ = sys.ForceDisable() sysproxy.ClearSentinel(cfgPath) fmt.Fprintf(m.stderr, "core exited: system proxy cleared\n") } return } } func (m *Manager) Disconnect() error { m.mu.Lock() sys := m.sys engine := m.engine m.engine = nil m.profile = nil m.watchGeneration++ // stop watcher m.unexpectedDrop = false if engine != nil { m.lastStop.Add(1) } m.mu.Unlock() var first error if sys != nil && (sys.Enabled() || sysproxy.HasSentinel(m.cfgPath)) { var err error if sys.Enabled() { err = sys.Disable() } else { err = sys.ForceDisable() } if err != nil { if err2 := sys.ForceDisable(); err2 != nil && first == nil { first = err2 } else if first == nil { first = err } } sysproxy.ClearSentinel(m.cfgPath) } if engine != nil { done := make(chan error, 1) go func() { defer m.lastStop.Done() done <- engine.Stop() }() select { case err := <-done: if err != nil && first == nil { first = err } case <-time.After(3 * time.Second): // Stop continues in background; Connect will wait on lastStop. if first == nil { first = fmt.Errorf("отключение заняло слишком много времени") } } } return first } func (m *Manager) Status() Status { m.mu.Lock() defer m.mu.Unlock() st := Status{SystemProxy: m.sys.Enabled()} if m.engine == nil || !m.engine.Running() { return st } st.Connected = true if m.profile != nil { st.Profile = m.profile.Name st.Protocol = m.profile.Protocol } if hp, ok := m.engine.LocalHTTPProxy(); ok { st.HTTPProxy = hp } if sp, ok := m.engine.LocalSOCKSProxy(); ok { st.SOCKSProxy = sp } return st } // Probe fetches url via the local HTTP or SOCKS proxy to verify the tunnel. func (m *Manager) Probe(ctx context.Context, testURL string) error { m.mu.Lock() engine := m.engine m.mu.Unlock() if engine == nil || !engine.Running() { return fmt.Errorf("not connected") } if testURL == "" { testURL = "https://www.google.com/generate_204" } var proxyURL *url.URL if hp, ok := engine.LocalHTTPProxy(); ok { proxyURL, _ = url.Parse("http://" + hp) } else if sp, ok := engine.LocalSOCKSProxy(); ok { proxyURL, _ = url.Parse("socks5://" + sp) } else { return fmt.Errorf("no local proxy listen address") } transport := &http.Transport{ Proxy: http.ProxyURL(proxyURL), DialContext: (&net.Dialer{ Timeout: 15 * time.Second, }).DialContext, TLSHandshakeTimeout: 15 * time.Second, ForceAttemptHTTP2: true, } client := &http.Client{ Transport: transport, Timeout: 30 * time.Second, CheckRedirect: func(req *http.Request, via []*http.Request) error { return http.ErrUseLastResponse }, } req, err := http.NewRequestWithContext(ctx, http.MethodGet, testURL, nil) if err != nil { return err } resp, err := client.Do(req) if err != nil { return fmt.Errorf("probe via %s: %w", proxyURL.String(), err) } defer resp.Body.Close() if resp.StatusCode >= 400 { return fmt.Errorf("probe unexpected status %d", resp.StatusCode) } return nil } func (m *Manager) BinDir() string { return m.binDir } func (m *Manager) ConfigPath() string { return m.cfgPath } // Config returns a deep snapshot for UI reads. Mutations must go through Manager methods. func (m *Manager) Config() *config.Config { m.mu.Lock() defer m.mu.Unlock() if m.cfg == nil { return &config.Config{} } return m.cfg.Clone() } // UIPoll is a cheap manager snapshot for GUI polling (no JSON deep-clone of all proxies). type UIPoll struct { Status Status Active string SystemProxy bool Subscription string ActiveProxy string ActiveProtocol config.Protocol Profiles []config.ProfileInfo Hy2 Hy2Options BinDir string LastError string ConnectOnLaunch bool AutoReconnect bool } // PollUI gathers status + profile list under one lock without Config().Clone(). func (m *Manager) PollUI(includeSecrets bool) UIPoll { m.mu.Lock() defer m.mu.Unlock() out := UIPoll{ Status: Status{SystemProxy: m.sys.Enabled()}, BinDir: m.binDir, SystemProxy: false, Subscription: "", Hy2: Hy2Options{Congestion: "bbr", BBRProfile: "standard"}, LastError: m.lastError, } if m.cfg != nil { out.Active = m.cfg.Active out.SystemProxy = m.cfg.SystemProxy out.Subscription = m.cfg.SubscriptionURL out.ConnectOnLaunch = m.cfg.ConnectOnLaunch out.AutoReconnect = m.cfg.AutoReconnect if includeSecrets { out.Profiles = m.cfg.ListProfiles() } else { out.Profiles = m.cfg.ListProfilesForPoll() } if p, err := m.cfg.ActiveProfile(); err == nil { out.Active = p.Name out.ActiveProxy = p.Proxy out.ActiveProtocol = p.Protocol out.Hy2 = Hy2Options{ Congestion: p.Hy2Congestion, BBRProfile: p.Hy2BBRProfile, BandwidthUp: p.Hy2BandwidthUp, BandwidthDown: p.Hy2BandwidthDown, Obfs: p.Hy2Obfs, ObfsPassword: p.Hy2ObfsPassword, SNI: p.Hy2SNI, Insecure: p.Hy2Insecure, PinSHA256: p.Hy2PinSHA256, FastOpen: p.Hy2FastOpen, Lazy: p.Hy2Lazy, HopInterval: p.Hy2HopInterval, } } } if m.engine != nil && m.engine.Running() { out.Status.Connected = true if m.profile != nil { out.Status.Profile = m.profile.Name out.Status.Protocol = m.profile.Protocol } if hp, ok := m.engine.LocalHTTPProxy(); ok { out.Status.HTTPProxy = hp } if sp, ok := m.engine.LocalSOCKSProxy(); ok { out.Status.SOCKSProxy = sp } } return out } func (m *Manager) SetSystemProxy(enabled bool) { m.mu.Lock() defer m.mu.Unlock() m.cfg.SystemProxy = enabled } // Prefs are GUI toggles persisted in config.json. type Prefs struct { ConnectOnLaunch bool `json:"connect_on_launch"` AutoReconnect bool `json:"auto_reconnect"` SystemProxy bool `json:"system_proxy"` } func (m *Manager) Prefs() Prefs { m.mu.Lock() defer m.mu.Unlock() if m.cfg == nil { return Prefs{} } return Prefs{ ConnectOnLaunch: m.cfg.ConnectOnLaunch, AutoReconnect: m.cfg.AutoReconnect, SystemProxy: m.cfg.SystemProxy, } } func (m *Manager) SavePrefs(p Prefs) error { m.mu.Lock() defer m.mu.Unlock() if m.cfg == nil { return fmt.Errorf("no config") } m.cfg.ConnectOnLaunch = p.ConnectOnLaunch m.cfg.AutoReconnect = p.AutoReconnect m.cfg.SystemProxy = p.SystemProxy return config.Save(m.cfgPath, *m.cfg) } func (m *Manager) UpdateProxyURI(proxy string) error { m.mu.Lock() defer m.mu.Unlock() active, err := m.cfg.ActiveProfile() if err != nil { return err } normalized, proto, _, err := linknorm.Normalize(active.Protocol, proxy) if err != nil { return err } if err := m.cfg.SetActiveProxy(normalized); err != nil { return err } if proto != "" { _ = m.cfg.UpsertProfileWithProtocol(m.cfg.Active, normalized, proto) } return config.Save(m.cfgPath, *m.cfg) } func (m *Manager) SetActiveProfile(name string) error { m.mu.Lock() defer m.mu.Unlock() if m.engine != nil && m.engine.Running() { return fmt.Errorf("сначала отключитесь") } if err := m.cfg.SetActive(name); err != nil { return err } return config.Save(m.cfgPath, *m.cfg) } func (m *Manager) SaveProfile(name, proxy string) error { m.mu.Lock() defer m.mu.Unlock() name = strings.TrimSpace(name) if name == "" { return fmt.Errorf("имя профиля пустое") } connected := m.engine != nil && m.engine.Running() if connected && name == m.cfg.Active { return fmt.Errorf("сначала отключитесь, чтобы изменить активный профиль") } proxy = strings.TrimSpace(proxy) var proto config.Protocol if proxy != "" { normalized, detected, _, err := linknorm.Normalize("", proxy) if err != nil { return err } proxy = normalized proto = detected } var err error if connected { err = m.cfg.UpsertProfileKeepActive(name, proxy, proto) } else { err = m.cfg.UpsertProfileWithProtocol(name, proxy, proto) } if err != nil { return err } if p, err := m.cfg.ActiveProfile(); err == nil { hysteria2.EnrichProfile(p) } return config.Save(m.cfgPath, *m.cfg) } // SaveActiveProfile updates the current profile (rename if needed) and proxy URI. func (m *Manager) SaveActiveProfile(name, proxy string) error { m.mu.Lock() defer m.mu.Unlock() if m.engine != nil && m.engine.Running() { return fmt.Errorf("сначала отключитесь") } name = strings.TrimSpace(name) if name == "" { return fmt.Errorf("имя профиля пустое") } proxy = strings.TrimSpace(proxy) var proto config.Protocol if proxy != "" { activeProto := config.Protocol("") if p, err := m.cfg.ActiveProfile(); err == nil { activeProto = p.Protocol } normalized, detected, _, err := linknorm.Normalize(activeProto, proxy) if err != nil { return err } proxy = normalized proto = detected } active := m.cfg.Active if name != active { exists := false for _, p := range m.cfg.Profiles { if p.Name == name { exists = true break } } if exists { if err := m.cfg.SetActive(name); err != nil { return err } } else if err := m.cfg.RenameProfile(active, name); err != nil { return err } } if err := m.cfg.UpsertProfileWithProtocol(name, proxy, proto); err != nil { return err } if p, err := m.cfg.ActiveProfile(); err == nil { hysteria2.EnrichProfile(p) } return config.Save(m.cfgPath, *m.cfg) } func (m *Manager) DeleteProfile(name string) error { m.mu.Lock() defer m.mu.Unlock() if m.engine != nil && m.engine.Running() { return fmt.Errorf("сначала отключитесь") } if err := m.cfg.DeleteProfile(name); err != nil { return err } return config.Save(m.cfgPath, *m.cfg) } func (m *Manager) Profiles() []config.ProfileInfo { m.mu.Lock() defer m.mu.Unlock() return m.cfg.ListProfiles() } // Hy2Options are editable Hysteria 2 extras for the active profile. type Hy2Options struct { Congestion string `json:"congestion"` BBRProfile string `json:"bbr_profile"` BandwidthUp string `json:"bandwidth_up"` BandwidthDown string `json:"bandwidth_down"` Obfs string `json:"obfs"` ObfsPassword string `json:"obfs_password"` SNI string `json:"sni"` Insecure bool `json:"insecure"` PinSHA256 string `json:"pin_sha256"` FastOpen bool `json:"fast_open"` Lazy bool `json:"lazy"` HopInterval string `json:"hop_interval"` } func (m *Manager) ActiveHy2Options() Hy2Options { m.mu.Lock() defer m.mu.Unlock() p, err := m.cfg.ActiveProfile() if err != nil { return Hy2Options{Congestion: "bbr", BBRProfile: "standard"} } return Hy2Options{ Congestion: p.Hy2Congestion, BBRProfile: p.Hy2BBRProfile, BandwidthUp: p.Hy2BandwidthUp, BandwidthDown: p.Hy2BandwidthDown, Obfs: p.Hy2Obfs, ObfsPassword: p.Hy2ObfsPassword, SNI: p.Hy2SNI, Insecure: p.Hy2Insecure, PinSHA256: p.Hy2PinSHA256, FastOpen: p.Hy2FastOpen, Lazy: p.Hy2Lazy, HopInterval: p.Hy2HopInterval, } } func (m *Manager) SaveHy2Options(opts Hy2Options) error { m.mu.Lock() defer m.mu.Unlock() if m.engine != nil && m.engine.Running() { return fmt.Errorf("сначала отключитесь") } p, err := m.cfg.ActiveProfile() if err != nil { return err } p.Protocol = config.ProtocolHysteria2 p.Hy2Congestion = strings.TrimSpace(opts.Congestion) p.Hy2BBRProfile = strings.TrimSpace(opts.BBRProfile) p.Hy2BandwidthUp = strings.TrimSpace(opts.BandwidthUp) p.Hy2BandwidthDown = strings.TrimSpace(opts.BandwidthDown) p.Hy2Obfs = strings.TrimSpace(opts.Obfs) p.Hy2ObfsPassword = strings.TrimSpace(opts.ObfsPassword) p.Hy2SNI = strings.TrimSpace(opts.SNI) p.Hy2Insecure = opts.Insecure p.Hy2PinSHA256 = strings.TrimSpace(opts.PinSHA256) p.Hy2FastOpen = opts.FastOpen p.Hy2Lazy = opts.Lazy p.Hy2HopInterval = strings.TrimSpace(opts.HopInterval) return config.Save(m.cfgPath, *m.cfg) } func (m *Manager) SubscriptionURL() string { m.mu.Lock() defer m.mu.Unlock() return m.cfg.SubscriptionURL } func (m *Manager) ImportSubscription(rawURL string) (int, error) { ctx, cancel := context.WithTimeout(context.Background(), 45*time.Second) defer cancel() items, err := subscription.Fetch(ctx, rawURL) if err != nil { return 0, err } if len(items) == 0 { return 0, fmt.Errorf("в подписке нет поддерживаемых ссылок (naive / hysteria2 / awg / vless / vmess / trojan)") } m.mu.Lock() defer m.mu.Unlock() if m.engine != nil && m.engine.Running() { return 0, fmt.Errorf("сначала отключитесь") } m.cfg.SubscriptionURL = strings.TrimSpace(rawURL) newNames := make([]string, 0, len(items)) seen := map[string]struct{}{} for _, it := range items { name := strings.TrimSpace(it.Name) if name == "" { continue } _ = m.cfg.UpsertProfileKeepActive(name, it.URI, it.Protocol) for i := range m.cfg.Profiles { if m.cfg.Profiles[i].Name == name { hysteria2.EnrichProfile(&m.cfg.Profiles[i]) break } } if _, ok := seen[name]; !ok { seen[name] = struct{}{} newNames = append(newNames, name) } } prev := map[string]struct{}{} for _, n := range m.cfg.SubscriptionNames { prev[n] = struct{}{} } // Prune profiles that came from the previous subscription sync but are gone now. if len(prev) > 0 { kept := make([]config.Profile, 0, len(m.cfg.Profiles)) for _, p := range m.cfg.Profiles { if _, wasSub := prev[p.Name]; wasSub { if _, still := seen[p.Name]; !still { continue } } kept = append(kept, p) } if len(kept) == 0 { // Safety: never wipe all profiles. kept = append([]config.Profile(nil), m.cfg.Profiles...) } else { m.cfg.Profiles = kept } activeOK := false for _, p := range m.cfg.Profiles { if p.Name == m.cfg.Active { activeOK = true break } } if !activeOK && len(m.cfg.Profiles) > 0 { m.cfg.Active = m.cfg.Profiles[0].Name } } m.cfg.SubscriptionNames = newNames if err := config.Save(m.cfgPath, *m.cfg); err != nil { return 0, err } return len(newNames), nil } func (m *Manager) SaveConfig() error { m.mu.Lock() defer m.mu.Unlock() return config.Save(m.cfgPath, *m.cfg) } func (m *Manager) Reload() error { m.mu.Lock() defer m.mu.Unlock() if m.engine != nil && m.engine.Running() { return fmt.Errorf("disconnect before reloading config") } cfg, err := config.Load(m.cfgPath) if err != nil { return err } binDir, err := config.ResolveBinDir(m.cfgPath, cfg.BinDir) if err != nil { return err } m.cfg = cfg m.binDir = binDir return nil } func (m *Manager) newEngine(p config.Protocol) (Engine, error) { switch p { case config.ProtocolNaive, "": return naive.New(m.stderr), nil case config.ProtocolHysteria2: return hysteria2.New(m.stderr), nil case config.ProtocolAWG: return awg.New(m.stderr), nil case config.ProtocolVLESS, config.ProtocolVMess, config.ProtocolTrojan: return xray.New(m.stderr), nil default: return nil, fmt.Errorf("unsupported protocol %q", p) } } // EnsureCore downloads the binary required by the active (or given) protocol. func (m *Manager) EnsureCore(proto config.Protocol) (string, error) { m.mu.Lock() binDir := m.binDir if proto == "" { if p, err := m.cfg.ActiveProfile(); err == nil { proto = p.Protocol } } m.mu.Unlock() var ( path string err error ) switch proto { case config.ProtocolHysteria2: path, err = hysteria2.EnsureBinary(binDir) case config.ProtocolAWG: path, err = awg.EnsureBinary(binDir) case config.ProtocolVLESS, config.ProtocolVMess, config.ProtocolTrojan: path, err = xray.EnsureBinary(binDir) default: path, err = naive.EnsureBinary(binDir) proto = config.ProtocolNaive } if err != nil { return "", err } corebin.Set(corebin.CacheKey(binDir, corebin.ProtoKey(string(proto))), path) return path, nil } // EnsureAllCores installs naive + hysteria2 + xray cores (AWG is embedded). // Downloads run in parallel (one goroutine per core) to use multiple CPU cores / network. func (m *Manager) EnsureAllCores() (map[string]string, error) { m.mu.Lock() binDir := m.binDir m.mu.Unlock() corebin.Invalidate() type item struct { key string path string err error } jobs := []struct { key string fn func() (string, error) }{ {"naive", func() (string, error) { return naive.EnsureBinary(binDir) }}, {"hysteria2", func() (string, error) { return hysteria2.EnsureBinary(binDir) }}, {"awg", func() (string, error) { return awg.EnsureBinary(binDir) }}, {"xray", func() (string, error) { return xray.EnsureBinary(binDir) }}, } ch := make(chan item, len(jobs)) for _, j := range jobs { j := j go func() { path, err := j.fn() ch <- item{key: j.key, path: path, err: err} }() } out := map[string]string{} var first error for range jobs { it := <-ch if it.err != nil { if first == nil { first = fmt.Errorf("%s: %w", it.key, it.err) } continue } out[it.key] = it.path corebin.Set(corebin.CacheKey(binDir, it.key), it.path) } return out, first }