package main import ( "context" "fmt" "log" "net" "os" "strconv" "time" "github.com/go-redis/redis/v9" "github.com/hashicorp/consul/api" "go.minekube.com/brigodier" . "go.minekube.com/common/minecraft/component" "go.minekube.com/common/minecraft/component/codec/legacy" "go.minekube.com/gate/cmd/gate" "go.minekube.com/gate/pkg/command" "go.minekube.com/gate/pkg/edition/java/proxy" "go.minekube.com/gate/pkg/runtime/event" ) func main() { // Add our "plug-in" to be initialized on Gate start. proxy.Plugins = append(proxy.Plugins, proxy.Plugin{ Name: "FenceProxy", Init: func(proxy *proxy.Proxy) error { return newFenceProxy(proxy).init() }, }) // Execute Gate entrypoint and block until shutdown. // We could also run gate.Start if we don't need Gate's command-line. gate.Execute() } // FenceProxy is a simple proxy that adds a `/broadcast` command // and sends a message on server switch. type FenceProxy struct { *proxy.Proxy legacyCodec *legacy.Legacy } func newFenceProxy(proxy *proxy.Proxy) *FenceProxy { return &FenceProxy{ Proxy: proxy, legacyCodec: &legacy.Legacy{Char: legacy.AmpersandChar}, } } // initialize our sample proxy func (p *FenceProxy) init() error { p.registerCommands() go p.registerServers() return p.registerSubscribers() } func (p *FenceProxy) registerServers() { // Get a new client config := api.DefaultConfig() config.Address = getEnv("CONSUL_ADDR", "127.0.0.1:8500") client, err := api.NewClient(config) if err != nil { panic(err) } // Get handle to catalog API catalog := client.Catalog() q := &api.QueryOptions{RequireConsistent: true} var lastServers []proxy.ServerInfo for { allSvc, _, err := catalog.Services(q) if err != nil { log.Printf("Error querying Consul: %v", err) } var minecraftServices []string for svc, tags := range allSvc { for _, tag := range tags { if tag == "consulity.enabled" { minecraftServices = append(minecraftServices, svc) break } } } var servers []proxy.ServerInfo for _, name := range minecraftServices { svcs, _, err := catalog.Service(name, "", q) if err != nil { log.Printf("Error querying Consul: %v", err) continue } for _, svc := range svcs { adr := fmt.Sprintf("%v:%d", svc.ServiceAddress, svc.ServicePort) ip, err := net.ResolveTCPAddr("tcp", adr) if err != nil { log.Printf("Error converting address: %s", adr) continue } servers = append(servers, proxy.NewServerInfo(svc.ID, ip)) } } // servers in lastServers that don't exist in servers now var notExistingServers []proxy.ServerInfo for _, server := range lastServers { stillExists := false for _, curServer := range servers { if curServer.Name() == server.Name() { stillExists = true } } if !stillExists { fmt.Printf("Server doesn't exist anymore: %v", server.Name()) notExistingServers = append(notExistingServers, server) } } for _, server := range notExistingServers { p.Unregister(server) } for _, server := range servers { alreadyExists := false for _, lastServer := range lastServers { if lastServer.Name() == server.Name() { alreadyExists = true } } if !alreadyExists { p.Register(server) } } lastServers = servers time.Sleep(5 * time.Second) } } // Register a proxy-wide commands (can be run while being on any server) func (p *FenceProxy) registerCommands() { // Registers the "/broadcast" command p.Command().Register(brigodier.Literal("broadcast").Then( // Adds message argument as in "/broadcast " brigodier.Argument("message", brigodier.StringPhrase). // Adds completion suggestions as in "/broadcast [suggestions]" Suggests(command.SuggestFunc(func( c *command.Context, b *brigodier.SuggestionsBuilder, ) *brigodier.Suggestions { player, ok := c.Source.(proxy.Player) if ok { b.Suggest("&oI am &6&l" + player.Username()) } b.Suggest("Hello world!") return b.Build() })). // Executed when running "/broadcast " Executes(command.Command(func(c *command.Context) error { // Colorize/format message message, err := p.legacyCodec.Unmarshal([]byte(c.String("message"))) if err != nil { return c.Source.SendMessage(&Text{ Content: fmt.Sprintf("Error formatting message: %v", err)}) } // Send to all players on this proxy for _, player := range p.Players() { // Send message in new goroutine, // to not halt loop on slow connections. go func(p proxy.Player) { _ = p.SendMessage(message) }(player) } return nil })), )) } var ctx = context.Background() var localPlayers int64 = 0 // Register event subscribers func (p *FenceProxy) registerSubscribers() error { rdb := redis.NewClient(&redis.Options{ Addr: getEnv("REDIS_ADDR", "localhost:6379"), Password: "", // no password set DB: 0, // use default DB }) p.Event().Subscribe(&proxy.PingEvent{}, 0, func(ev event.Event) { e := ev.(*proxy.PingEvent) p := e.Ping() val, err := rdb.Get(ctx, "players").Result() if err == redis.Nil { // does not exist yet rdb.Set(ctx, "players", 0, 0) } else if err != nil { panic(err) } p.Players.Online, _ = strconv.Atoi(val) p.Players.Max = p.Players.Online + 1 }) p.Event().Subscribe(&proxy.LoginEvent{}, 0, func(ev event.Event) { //e := ev.(*proxy.LoginEvent) //e.Deny(&Text{Content: "&cMaintenance mode"}) }) p.Event().Subscribe(&proxy.PlayerChooseInitialServerEvent{}, 0, func(ev event.Event) { e := ev.(*proxy.PlayerChooseInitialServerEvent) e.SetInitialServer(p.Servers()[0]) localPlayers += 1 rdb.Incr(ctx, "players") }) p.Event().Subscribe(&proxy.DisconnectEvent{}, 0, func(ev event.Event) { localPlayers -= 1 rdb.Decr(ctx, "players") }) p.Event().Subscribe(&proxy.ShutdownEvent{}, 0, func(ev event.Event) { rdb.DecrBy(ctx, "players", localPlayers) }) return nil } func getEnv(key, fallback string) string { if value, ok := os.LookupEnv(key); ok { return value } return fallback }