mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-09-30 06:43:48 +02:00
Update to Go 1.26.4, build runners and go.mod depedencies (#3108)
* update golang and dependencies * fix incorrect log formatting * clean mod chache and introduce GOARCH in Dockerfile (choose dynamically) * remove GO111MODULE mentions * bump github actions from v2 to v3 * bump docker runners to v7 * use extra event store for backwards compatibility with existing codebase * updated generated opconfig api
This commit is contained in:
@@ -65,6 +65,7 @@ type Controller struct {
|
||||
nodesInformer cache.SharedIndexInformer
|
||||
podCh chan cluster.PodEvent
|
||||
|
||||
clusterEventStores []cache.Store // [workerID]Store
|
||||
clusterEventQueues []*cache.FIFO // [workerID]Queue
|
||||
lastClusterSyncTime int64
|
||||
lastClusterRepairTime int64
|
||||
@@ -356,17 +357,19 @@ func (c *Controller) initController() {
|
||||
c.config.InfrastructureRoles = infraRoles
|
||||
}
|
||||
|
||||
c.clusterEventStores = make([]cache.Store, c.opConfig.Workers)
|
||||
c.clusterEventQueues = make([]*cache.FIFO, c.opConfig.Workers)
|
||||
c.workerLogs = make(map[uint32]ringlog.RingLogger, c.opConfig.Workers)
|
||||
for i := range c.clusterEventQueues {
|
||||
c.clusterEventQueues[i] = cache.NewFIFO(func(obj interface{}) (string, error) {
|
||||
keyFn := func(obj interface{}) (string, error) {
|
||||
e, ok := obj.(ClusterEvent)
|
||||
if !ok {
|
||||
return "", fmt.Errorf("could not cast to ClusterEvent")
|
||||
return "", fmt.Errorf("could not cast to cluster event")
|
||||
}
|
||||
|
||||
return queueClusterKey(e.EventType, e.UID), nil
|
||||
})
|
||||
}
|
||||
c.clusterEventStores[i] = cache.NewStore(keyFn)
|
||||
c.clusterEventQueues[i] = cache.NewFIFO(keyFn)
|
||||
}
|
||||
|
||||
c.apiserver = apiserver.New(c, c.opConfig.APIPort, c.logger.Logger)
|
||||
|
||||
@@ -80,8 +80,8 @@ func (c *Controller) GetStatus() *spec.ControllerStatus {
|
||||
c.clustersMu.RUnlock()
|
||||
|
||||
queueSizes := make(map[int]int, c.opConfig.Workers)
|
||||
for workerID, queue := range c.clusterEventQueues {
|
||||
queueSizes[workerID] = len(queue.ListKeys())
|
||||
for workerID, store := range c.clusterEventStores {
|
||||
queueSizes[workerID] = len(store.ListKeys())
|
||||
}
|
||||
|
||||
return &spec.ControllerStatus{
|
||||
@@ -180,11 +180,11 @@ func (c *Controller) Fire(e *logrus.Entry) error {
|
||||
|
||||
// ListQueue dumps cluster event queue of the provided worker
|
||||
func (c *Controller) ListQueue(workerID uint32) (*spec.QueueDump, error) {
|
||||
if workerID >= uint32(len(c.clusterEventQueues)) {
|
||||
if workerID >= uint32(len(c.clusterEventStores)) {
|
||||
return nil, fmt.Errorf("could not find worker")
|
||||
}
|
||||
|
||||
q := c.clusterEventQueues[workerID]
|
||||
q := c.clusterEventStores[workerID]
|
||||
return &spec.QueueDump{
|
||||
Keys: q.ListKeys(),
|
||||
List: q.List(),
|
||||
@@ -196,7 +196,7 @@ func (c *Controller) GetWorkersCnt() uint32 {
|
||||
return c.opConfig.Workers
|
||||
}
|
||||
|
||||
//WorkerStatus provides status of the worker
|
||||
// WorkerStatus provides status of the worker
|
||||
func (c *Controller) WorkerStatus(workerID uint32) (*cluster.WorkerStatus, error) {
|
||||
obj, ok := c.curWorkerCluster.Load(workerID)
|
||||
if !ok || obj == nil {
|
||||
|
||||
@@ -182,7 +182,7 @@ func (c *Controller) addCluster(lg *logrus.Entry, clusterName spec.NamespacedNam
|
||||
return cl, nil
|
||||
}
|
||||
|
||||
func (c *Controller) processEvent(event ClusterEvent) {
|
||||
func (c *Controller) processEvent(event ClusterEvent, isInInitialList bool) {
|
||||
var clusterName spec.NamespacedName
|
||||
var clHistory ringlog.RingLogger
|
||||
var err error
|
||||
@@ -371,11 +371,22 @@ func (c *Controller) processClusterEventsQueue(idx int, stopCh <-chan struct{},
|
||||
|
||||
go func() {
|
||||
<-stopCh
|
||||
c.clusterEventQueues[idx].Close()
|
||||
(*c.clusterEventQueues[idx]).Close()
|
||||
}()
|
||||
|
||||
for {
|
||||
obj, err := c.clusterEventQueues[idx].Pop(cache.PopProcessFunc(func(interface{}, bool) error { return nil }))
|
||||
_, err := (*c.clusterEventQueues[idx]).Pop(cache.PopProcessFunc(func(obj interface{}, isInitialList bool) error {
|
||||
event, ok := obj.(ClusterEvent)
|
||||
if !ok {
|
||||
c.logger.Errorf("could not cast to cluster event")
|
||||
return nil // skip event to keep processing
|
||||
}
|
||||
c.processEvent(event, isInitialList)
|
||||
if err := c.clusterEventStores[idx].Delete(obj); err != nil {
|
||||
c.logger.Errorf("failed to delete key from lookup store: %v", err)
|
||||
}
|
||||
return nil
|
||||
}))
|
||||
if err != nil {
|
||||
if err == cache.ErrFIFOClosed {
|
||||
return
|
||||
@@ -383,12 +394,6 @@ func (c *Controller) processClusterEventsQueue(idx int, stopCh <-chan struct{},
|
||||
c.logger.Errorf("error when processing cluster events queue: %v", err)
|
||||
continue
|
||||
}
|
||||
event, ok := obj.(ClusterEvent)
|
||||
if !ok {
|
||||
c.logger.Errorf("could not cast to ClusterEvent")
|
||||
}
|
||||
|
||||
c.processEvent(event)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -523,7 +528,10 @@ func (c *Controller) queueClusterEvent(informerOldSpec, informerNewSpec *acidv1.
|
||||
}
|
||||
|
||||
lg := c.logger.WithField("worker", workerID).WithField("cluster-name", clusterName)
|
||||
if err := c.clusterEventQueues[workerID].Add(clusterEvent); err != nil {
|
||||
if err := c.clusterEventStores[workerID].Add(clusterEvent); err != nil {
|
||||
lg.Errorf("error while storing cluster event for lookup: %v", clusterEvent)
|
||||
}
|
||||
if err := (*c.clusterEventQueues[workerID]).Add(clusterEvent); err != nil {
|
||||
lg.Errorf("error while queueing cluster event: %v", clusterEvent)
|
||||
}
|
||||
lg.Infof("%s event has been queued", eventType)
|
||||
@@ -533,9 +541,9 @@ func (c *Controller) queueClusterEvent(informerOldSpec, informerNewSpec *acidv1.
|
||||
}
|
||||
// A delete event discards all prior requests for that cluster.
|
||||
for _, evType := range []EventType{EventAdd, EventSync, EventUpdate, EventRepair} {
|
||||
obj, exists, err := c.clusterEventQueues[workerID].GetByKey(queueClusterKey(evType, uid))
|
||||
obj, exists, err := c.clusterEventStores[workerID].GetByKey(queueClusterKey(evType, uid))
|
||||
if err != nil {
|
||||
lg.Warningf("could not get event from the queue: %v", err)
|
||||
lg.Warningf("could not get event from the lookup store: %v", err)
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -543,12 +551,18 @@ func (c *Controller) queueClusterEvent(informerOldSpec, informerNewSpec *acidv1.
|
||||
continue
|
||||
}
|
||||
|
||||
err = c.clusterEventQueues[workerID].Delete(obj)
|
||||
err = (*c.clusterEventQueues[workerID]).Delete(obj)
|
||||
if err != nil {
|
||||
lg.Warningf("could not delete event from the queue: %v", err)
|
||||
} else {
|
||||
lg.Debugf("event %s has been discarded for the cluster", evType)
|
||||
}
|
||||
err = c.clusterEventStores[workerID].Delete(obj)
|
||||
if err != nil {
|
||||
lg.Warningf("could not delete event from the lookup store: %v", err)
|
||||
} else {
|
||||
lg.Debugf("event %s has been deleted from the lookup store", evType)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -16,7 +16,7 @@ import (
|
||||
"github.com/zalando/postgres-operator/pkg/util"
|
||||
"github.com/zalando/postgres-operator/pkg/util/config"
|
||||
"github.com/zalando/postgres-operator/pkg/util/k8sutil"
|
||||
"gopkg.in/yaml.v2"
|
||||
"gopkg.in/yaml.v3"
|
||||
)
|
||||
|
||||
func (c *Controller) makeClusterConfig() cluster.Config {
|
||||
|
||||
Reference in New Issue
Block a user