package proxy import ( "context" "fmt" "io" "log/slog" "net" "time" "github.com/hexahost/gamecloud/edge-gateway/internal/config" "github.com/hexahost/gamecloud/edge-gateway/internal/controlplane" "github.com/hexahost/gamecloud/edge-gateway/internal/mc" ) type Gateway struct { cfg config.Config client *controlplane.Client log *slog.Logger } func New(cfg config.Config, log *slog.Logger) *Gateway { return &Gateway{ cfg: cfg, client: controlplane.NewClient(cfg.APIURL, cfg.EdgeKey), log: log, } } func (g *Gateway) ListenAndServe(ctx context.Context) error { listener, err := net.Listen("tcp", g.cfg.ListenAddr) if err != nil { return err } defer listener.Close() g.log.Info("edge gateway listening", "addr", g.cfg.ListenAddr, "playDomain", g.cfg.PlayDomain, "api", g.cfg.APIURL, ) go func() { <-ctx.Done() _ = listener.Close() }() for { conn, err := listener.Accept() if err != nil { select { case <-ctx.Done(): return nil default: g.log.Error("accept failed", "error", err) continue } } go g.handleConnection(conn) } } func (g *Gateway) handleConnection(clientConn net.Conn) { defer clientConn.Close() clientIP := clientConn.RemoteAddr().String() if host, _, err := net.SplitHostPort(clientIP); err == nil { clientIP = host } handshake, err := mc.ReadHandshake(clientConn) if err != nil { g.log.Warn("handshake failed", "error", err, "client", clientIP) return } slug := mc.ExtractSlug(handshake.ServerAddress, g.cfg.PlayDomain) if slug == "" { g.log.Warn("unknown hostname", "address", handshake.ServerAddress, "client", clientIP) return } resolve, err := g.client.Resolve(slug, clientIP) if err != nil { g.log.Warn("resolve failed", "slug", slug, "error", err) return } switch resolve.Action { case "start": startResp, err := g.client.Start(slug, clientIP) if err != nil { g.log.Warn("join start failed", "slug", slug, "error", err) return } resolve = &startResp.Resolve case "wait": resolve, err = g.waitForProxy(slug, clientIP) if err != nil { g.log.Warn("wait for server failed", "slug", slug, "error", err) return } case "reject": g.log.Info("connection rejected", "slug", slug, "message", resolve.Message) return } if resolve.Action != "proxy" || resolve.Backend == nil { g.log.Warn("no backend available", "slug", slug, "action", resolve.Action) return } backendAddr := fmt.Sprintf("%s:%d", resolve.Backend.Host, resolve.Backend.Port) backendConn, err := net.DialTimeout("tcp", backendAddr, 10*time.Second) if err != nil { g.log.Error("backend dial failed", "slug", slug, "backend", backendAddr, "error", err) return } defer backendConn.Close() if err := mc.WritePacket(backendConn, handshake.Raw); err != nil { g.log.Error("replay handshake failed", "slug", slug, "error", err) return } g.log.Info("proxying connection", "slug", slug, "client", clientIP, "backend", backendAddr, "status", resolve.Status, ) pipe(clientConn, backendConn) } func (g *Gateway) waitForProxy(slug, clientIP string) (*controlplane.ResolveResponse, error) { deadline := time.Now().Add(g.cfg.StartWait) for time.Now().Before(deadline) { resolve, err := g.client.Resolve(slug, clientIP) if err != nil { return nil, err } if resolve.Action == "proxy" { return resolve, nil } if resolve.Action == "reject" { return resolve, fmt.Errorf(resolve.Message) } time.Sleep(g.cfg.PollInterval) } return nil, fmt.Errorf("timed out waiting for server to start") } func pipe(left, right net.Conn) { errCh := make(chan error, 2) go func() { _, err := io.Copy(right, left); errCh <- err }() go func() { _, err := io.Copy(left, right); errCh <- err }() <-errCh }