package handler import ( "bytes" "context" "crypto/sha256" "encoding/json" "fmt" "io" "net" "net/http" "net/url" "os" "os/exec" "path/filepath" "strings" "time" "github.com/gofiber/fiber/v3" "github.com/peterqiu0516/sub-store/internal/model" "github.com/peterqiu0516/sub-store/internal/render" "github.com/peterqiu0516/sub-store/internal/util" ) func (d *Deps) HandleEgressInfo(c fiber.Ctx) error { var node model.ProxyNode if err := json.Unmarshal(c.Body(), &node); err != nil { return failed(c, "Invalid JSON", 400) } if getStringValue(node["server"]) == "" { node["server"] = getStringValue(node["address"]) } if getStringValue(node["name"]) == "" { node["name"] = "PROXY" } cacheKey := egressCacheKey(node) if entry, ok := d.CacheRepo.SafeGet(cacheKey); ok { var cached map[string]any if json.Unmarshal([]byte(entry.Content), &cached) == nil { cached["cached"] = true return success(c, cached) } } latencyMs, latencyErr := probeServerPortLatency(node, 5*time.Second) port, err := freeLocalPort() if err != nil { return failed(c, err.Error(), 500) } configData, err := buildEgressProbeConfig(node, port) if err != nil { return failed(c, err.Error(), 400) } info, err := runEgressProbe(configData, port) if err != nil { info = fiber.Map{"egressError": err.Error()} } if latencyMs >= 0 { info["latencyMs"] = latencyMs } if latencyErr != "" { info["latencyError"] = latencyErr } info["cached"] = false if data, err := json.Marshal(info); err == nil { ttl := int(d.Cfg.Fetcher.CacheTTL.Seconds()) if ttl <= 0 { ttl = 300 } d.CacheRepo.SafePut(cacheKey, string(data), nil, ttl) } return success(c, info) } func egressCacheKey(node model.ProxyNode) string { clean := model.ProxyNode{} skip := map[string]bool{ "id": true, "name": true, "latencyMs": true, "latencyError": true, "egressIp": true, "egressCountry": true, "egressRegion": true, "egressError": true, "country": true, "region": true, "city": true, "isp": 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)) if !ok { continue } var cached map[string]any if json.Unmarshal([]byte(entry.Content), &cached) != nil { continue } for _, key := range []string{"egressIp", "country", "region", "city", "isp", "latencyMs", "latencyError", "egressError"} { if v, ok := cached[key]; ok { node[key] = v } } node["cached"] = true } return nodes } func probeServerPortLatency(node model.ProxyNode, timeout time.Duration) (int64, string) { server := getStringValue(node["server"]) port := toIntSafe(node["port"]) if server == "" || port <= 0 { return -1, "missing server or port" } start := time.Now() conn, err := net.DialTimeout("tcp", net.JoinHostPort(server, fmt.Sprint(port)), timeout) if err != nil { return -1, err.Error() } _ = conn.Close() return time.Since(start).Milliseconds(), "" } func buildEgressProbeConfig(node model.ProxyNode, port int) ([]byte, error) { probeNode := model.ProxyNode{} for k, v := range node { probeNode[k] = v } probeNode["name"] = "PROXY" outbound := render.ToSingBoxOutbound(probeNode) if outbound == nil { return nil, fmt.Errorf("Unsupported proxy node for egress probe") } doc := map[string]any{ "log": map[string]any{"level": "warn"}, "inbounds": []any{ map[string]any{ "type": "mixed", "tag": "mixed-in", "listen": "127.0.0.1", "listen_port": port, }, }, "outbounds": []any{ outbound, map[string]any{"type": "direct", "tag": "DIRECT"}, }, "route": map[string]any{ "auto_detect_interface": true, "final": "PROXY", "rules": []any{map[string]any{"action": "sniff"}}, }, } return json.Marshal(doc) } func runEgressProbe(configData []byte, port int) (fiber.Map, error) { singBox, err := exec.LookPath("sing-box") if err != nil { return nil, fmt.Errorf("sing-box executable not found") } dir, err := os.MkdirTemp("", "sub-store-egress-*") if err != nil { return nil, err } defer os.RemoveAll(dir) configPath := filepath.Join(dir, "config.json") if err := os.WriteFile(configPath, configData, 0o600); err != nil { return nil, err } ctx, cancel := context.WithTimeout(context.Background(), 25*time.Second) defer cancel() cmd := exec.CommandContext(ctx, singBox, "run", "-c", configPath) var stderr bytes.Buffer cmd.Stderr = &stderr if err := cmd.Start(); err != nil { return nil, err } defer func() { cancel() _ = cmd.Wait() }() if err := waitTCP("127.0.0.1", port, 5*time.Second); err != nil { msg := strings.TrimSpace(stderr.String()) if msg != "" { return nil, fmt.Errorf("%s", msg) } return nil, err } proxyURL, _ := url.Parse(fmt.Sprintf("http://127.0.0.1:%d", port)) client := &http.Client{ Timeout: 15 * time.Second, Transport: &http.Transport{ Proxy: http.ProxyURL(proxyURL), }, } resp, err := client.Get("https://ipwho.is/?lang=en") if err != nil { return nil, err } defer resp.Body.Close() body, _ := io.ReadAll(io.LimitReader(resp.Body, util.MaxFlowRespBytes)) var data map[string]any if err := json.Unmarshal(body, &data); err != nil { return nil, fmt.Errorf("Invalid egress info response") } if success, ok := data["success"].(bool); ok && !success { msg := getStringValue(data["message"]) if msg == "" { msg = "Egress info lookup failed" } return nil, fmt.Errorf("%s", msg) } connection, _ := data["connection"].(map[string]any) return fiber.Map{ "egressIp": data["ip"], "country": data["country"], "region": data["region"], "city": data["city"], "isp": connection["isp"], }, nil } func freeLocalPort() (int, error) { l, err := net.Listen("tcp", "127.0.0.1:0") if err != nil { return 0, err } defer l.Close() return l.Addr().(*net.TCPAddr).Port, nil } func waitTCP(host string, port int, timeout time.Duration) error { deadline := time.Now().Add(timeout) addr := fmt.Sprintf("%s:%d", host, port) for time.Now().Before(deadline) { conn, err := net.DialTimeout("tcp", addr, 200*time.Millisecond) if err == nil { conn.Close() return nil } time.Sleep(100 * time.Millisecond) } return fmt.Errorf("Timed out waiting for sing-box") }