summaryrefslogtreecommitdiffstats
path: root/internal/nymclient
diff options
context:
space:
mode:
Diffstat (limited to 'internal/nymclient')
-rw-r--r--internal/nymclient/manager.go176
1 files changed, 176 insertions, 0 deletions
diff --git a/internal/nymclient/manager.go b/internal/nymclient/manager.go
new file mode 100644
index 0000000..3c0b15d
--- /dev/null
+++ b/internal/nymclient/manager.go
@@ -0,0 +1,176 @@
+package nymclient
+
+import (
+ "context"
+ "fmt"
+ "io"
+ "net"
+ "os"
+ "os/exec"
+ "path/filepath"
+ "strings"
+ "sync"
+ "time"
+)
+
+type Config struct {
+ Binary string
+ Provider string
+ HomeDir string
+ ClientID string
+ SocksAddr string
+ AnonymousReplies bool
+ StartupTimeout time.Duration
+ LogOutput io.Writer
+}
+
+type Manager struct {
+ cfg Config
+ cmd *exec.Cmd
+ mu sync.Mutex
+}
+
+func New(cfg Config) (*Manager, error) {
+ if cfg.Binary == "" {
+ cfg.Binary = "nym-socks5-client"
+ }
+ if cfg.HomeDir == "" {
+ return nil, fmt.Errorf("nym home dir is required")
+ }
+ if cfg.ClientID == "" {
+ cfg.ClientID = "n2usenet"
+ }
+ if cfg.SocksAddr == "" {
+ cfg.SocksAddr = "127.0.0.1:11080"
+ }
+ if cfg.Provider == "" {
+ return nil, fmt.Errorf("nym provider is required")
+ }
+ if cfg.StartupTimeout <= 0 {
+ cfg.StartupTimeout = 120 * time.Second
+ }
+ if cfg.LogOutput == nil {
+ cfg.LogOutput = os.Stderr
+ }
+ return &Manager{cfg: cfg}, nil
+}
+
+func (m *Manager) Init(ctx context.Context) error {
+ binary, err := exec.LookPath(m.cfg.Binary)
+ if err != nil {
+ if _, statErr := os.Stat(m.cfg.Binary); statErr != nil {
+ return fmt.Errorf("find nym-socks5-client: %w", err)
+ }
+ binary = m.cfg.Binary
+ }
+ m.cfg.Binary = binary
+
+ if err := os.MkdirAll(m.cfg.HomeDir, 0700); err != nil {
+ return fmt.Errorf("create nym home: %w", err)
+ }
+
+ configFile := filepath.Join(m.cfg.HomeDir, ".nym", "socks5-clients", m.cfg.ClientID, "config", "config.toml")
+ if data, err := os.ReadFile(configFile); err == nil {
+ if !strings.Contains(string(data), m.cfg.Provider) {
+ return fmt.Errorf("existing nym client config uses a different provider; remove %s to reinitialize intentionally", filepath.Dir(filepath.Dir(configFile)))
+ }
+ return nil
+ }
+
+ host, port, err := net.SplitHostPort(m.cfg.SocksAddr)
+ if err != nil {
+ return fmt.Errorf("split socks addr: %w", err)
+ }
+
+ args := []string{
+ "init",
+ "--id", m.cfg.ClientID,
+ "--provider", m.cfg.Provider,
+ "--host", host,
+ "--port", port,
+ }
+ if m.cfg.AnonymousReplies {
+ args = append(args, "--use-reply-surbs", "true")
+ }
+
+ cmd := exec.CommandContext(ctx, m.cfg.Binary, args...)
+ cmd.Env = append(os.Environ(), "HOME="+m.cfg.HomeDir)
+ cmd.Stdout = m.cfg.LogOutput
+ cmd.Stderr = m.cfg.LogOutput
+ if err := cmd.Run(); err != nil {
+ return fmt.Errorf("nym-socks5-client init: %w", err)
+ }
+ return nil
+}
+
+func (m *Manager) Start(ctx context.Context) error {
+ m.mu.Lock()
+ defer m.mu.Unlock()
+ if m.cmd != nil {
+ return nil
+ }
+
+ host, port, err := net.SplitHostPort(m.cfg.SocksAddr)
+ if err != nil {
+ return fmt.Errorf("split socks addr: %w", err)
+ }
+
+ args := []string{"run", "--id", m.cfg.ClientID, "--host", host, "--port", port}
+ if m.cfg.AnonymousReplies {
+ args = append(args, "--use-anonymous-replies", "true")
+ }
+
+ cmd := exec.CommandContext(ctx, m.cfg.Binary, args...)
+ cmd.Env = append(os.Environ(), "HOME="+m.cfg.HomeDir)
+ cmd.Stdout = m.cfg.LogOutput
+ cmd.Stderr = m.cfg.LogOutput
+ if err := cmd.Start(); err != nil {
+ return fmt.Errorf("start nym-socks5-client: %w", err)
+ }
+ m.cmd = cmd
+
+ go func() {
+ _ = cmd.Wait()
+ }()
+
+ deadline := time.Now().Add(m.cfg.StartupTimeout)
+ for {
+ dialCtx, cancel := context.WithTimeout(ctx, 500*time.Millisecond)
+ conn, err := (&net.Dialer{}).DialContext(dialCtx, "tcp", m.cfg.SocksAddr)
+ cancel()
+ if err == nil {
+ _ = conn.Close()
+ return nil
+ }
+ if time.Now().After(deadline) {
+ m.Stop()
+ return fmt.Errorf("nym-socks5-client did not open %s within %s", m.cfg.SocksAddr, m.cfg.StartupTimeout)
+ }
+ select {
+ case <-ctx.Done():
+ m.Stop()
+ return ctx.Err()
+ case <-time.After(500 * time.Millisecond):
+ }
+ }
+}
+
+func (m *Manager) Stop() {
+ m.mu.Lock()
+ defer m.mu.Unlock()
+ if m.cmd == nil || m.cmd.Process == nil {
+ return
+ }
+ _ = m.cmd.Process.Signal(os.Interrupt)
+ done := make(chan struct{})
+ go func() {
+ _ = m.cmd.Wait()
+ close(done)
+ }()
+ select {
+ case <-done:
+ case <-time.After(5 * time.Second):
+ _ = m.cmd.Process.Kill()
+ }
+ m.cmd = nil
+}