package controller import ( "context" "os/exec" "time" "git.giftfish.de/ston1th/haproxy-kb/pkg/cluster" "git.giftfish.de/ston1th/haproxy-kb/pkg/config" "git.giftfish.de/ston1th/haproxy-kb/pkg/haproxy" "git.giftfish.de/ston1th/haproxy-kb/pkg/vip" "github.com/go-logr/logr" ) type UpdateAction int const ( Add UpdateAction = iota Delete ) type ConfigUpdate struct { Action UpdateAction Config haproxy.LoadBalancer } func readConfigUpdates(c chan ConfigUpdate) (cu []ConfigUpdate) { for len(c) > 0 { select { case config := <-c: cu = append(cu, config) } } return } func applyConfigUpdates(cc cluster.CallbackContext, cu []ConfigUpdate) error { for _, c := range cu { switch c.Action { case Add: addIPs(cc, n, []string{c.Config.IP}) cc.Set(c.Config.IP, nil) case Delete: cc.Delete(c.Config.IP) deleteIPs(cc, n, []string{c.Config.IP}) } } } func NewLBController(cfg *config.Config) (callbacks cluster.Callbacks, err error) { n, err := vip.NewNetworkWithLabel(cfg.Interface, cfg.Label) if err != nil { return } //TODO ha, err := haproxy.NewHAProxyManager("") if err != nil { return } configUpdate := make(chan ConfigUpdate, 100) callbacks = cluster.Callbacks{ Leader: func(ctx context.Context, cc cluster.CallbackContext) { t := time.NewTicker(time.Second * 10) if cfg.LeaderHook != "" { go func() { err := exec.CommandContext(ctx, cfg.LeaderHook).Run() if err != nil { cc.Error(err, "error running hook", "leaderHook", cfg.LeaderHook) } }() } for { cc.Info("leading", "id", cc.ID()) cfg, err := cc.Get("config") if err != nil { cc.Error("error getting config key", err) <-t.C continue } addIPs(cc, n, cfg.VirtualIPs) var lbs haproxy.Config err = ha.UpdateConfig(ctx, lbs) if err != nil { cc.Error("error updating haproxy config", err) <-t.C continue } select { case <-ctx.Done(): t.Stop() deleteIPs(cc, n, cfg.VirtualIPs) cc.Info("leading canceled", "id", cc.ID()) return case <-t.C: } } }, Follower: func(ctx context.Context, cc cluster.CallbackContext) { t := time.NewTicker(time.Second * 10) if cfg.FollowerHook != "" { go func() { err := exec.CommandContext(ctx, cfg.FollowerHook).Run() if err != nil { cc.Error(err, "error running hook", "followerHook", cfg.FollowerHook) } }() } for { cc.Info("following", "id", cc.ID()) //TODO //deleteIPs(cc, n, cfg.VirtualIPs) select { case <-ctx.Done(): t.Stop() cc.Info("following canceled", "id", cc.ID()) return case <-t.C: } } }, Cleanup: func(ctx context.Context, cc cluster.CallbackContext) { deleteIPs(cc, n, cfg.VirtualIPs) }, } return } func addIPs(log logr.Logger, n *vip.Network, cidrs []string) error { for _, cidr := range cidrs { err := n.AddIP(cidr) if err != nil { log.Error(err, "adding ip", "interface", n.Interface(), "ip", cidr) } } return nil } func deleteIPs(log logr.Logger, n *vip.Network, cidrs []string) error { for _, cidr := range cidrs { err := n.DeleteIP(cidr) if err != nil { log.Error(err, "deleting ip", "interface", n.Interface(), "ip", cidr) } } return nil }