mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-09-30 18:33:02 +02:00
Merge branch 'master' into support-many-namespaces
This commit is contained in:
@@ -101,8 +101,10 @@ func (c *Controller) initOperatorConfig() {
|
||||
watchedNsEnvVar, isPresentInOperatorEnv := os.LookupEnv("WATCHED_NAMESPACE")
|
||||
|
||||
if (!isPresentInOperatorConfigMap) && (!isPresentInOperatorEnv) {
|
||||
c.logger.Infoln("Neither the operator config map nor operator pod's environment define a namespace to watch. Fall back to watching the 'default' namespace.")
|
||||
configMapData["watched_namespace"] = v1.NamespaceDefault
|
||||
|
||||
c.logger.Infof("No namespace to watch specified. By convention, the operator falls back to watching the namespace it is deployed to: '%v' \n", spec.GetOperatorNamespace())
|
||||
configMapData["watched_namespace"] = spec.GetOperatorNamespace()
|
||||
|
||||
}
|
||||
|
||||
if (isPresentInOperatorConfigMap) && (!isPresentInOperatorEnv) {
|
||||
|
||||
+25
-3
@@ -3,6 +3,8 @@ package spec
|
||||
import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -26,6 +28,8 @@ const (
|
||||
EventUpdate EventType = "UPDATE"
|
||||
EventDelete EventType = "DELETE"
|
||||
EventSync EventType = "SYNC"
|
||||
|
||||
fileWithNamespace = "/var/run/secrets/kubernetes.io/serviceaccount/namespace"
|
||||
)
|
||||
|
||||
// ClusterEvent carries the payload of the Cluster TPR events.
|
||||
@@ -161,20 +165,38 @@ func (n NamespacedName) MarshalJSON() ([]byte, error) {
|
||||
|
||||
// Decode converts a (possibly unqualified) string into the namespaced name object.
|
||||
func (n *NamespacedName) Decode(value string) error {
|
||||
return n.DecodeWorker(value, GetOperatorNamespace())
|
||||
}
|
||||
|
||||
// DecodeWorker separates the decode logic to (unit) test
|
||||
// from obtaining the operator namespace that depends on k8s mounting files at runtime
|
||||
func (n *NamespacedName) DecodeWorker(value, operatorNamespace string) error {
|
||||
name := types.NewNamespacedNameFromString(value)
|
||||
|
||||
if strings.Trim(value, string(types.Separator)) != "" && name == (types.NamespacedName{}) {
|
||||
name.Name = value
|
||||
name.Namespace = v1.NamespaceDefault
|
||||
name.Namespace = operatorNamespace
|
||||
} else if name.Namespace == "" {
|
||||
name.Namespace = v1.NamespaceDefault
|
||||
name.Namespace = operatorNamespace
|
||||
}
|
||||
|
||||
if name.Name == "" {
|
||||
return fmt.Errorf("incorrect namespaced name")
|
||||
return fmt.Errorf("incorrect namespaced name: %v", value)
|
||||
}
|
||||
|
||||
*n = NamespacedName(name)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetOperatorNamespace assumes serviceaccount secret is mounted by kubernetes
|
||||
// Placing this func here instead of pgk/util avoids circular import
|
||||
func GetOperatorNamespace() string {
|
||||
|
||||
operatorNamespaceBytes, err := ioutil.ReadFile(fileWithNamespace)
|
||||
if err != nil {
|
||||
log.Fatalf("Unable to detect operator namespace from within its pod due to: %v", err)
|
||||
}
|
||||
|
||||
return string(operatorNamespaceBytes)
|
||||
}
|
||||
|
||||
+11
-5
@@ -5,22 +5,27 @@ import (
|
||||
"testing"
|
||||
)
|
||||
|
||||
const (
|
||||
mockOperatorNamespace = "acid"
|
||||
)
|
||||
|
||||
var nnTests = []struct {
|
||||
s string
|
||||
expected NamespacedName
|
||||
expectedMarshal []byte
|
||||
}{
|
||||
{`acid/cluster`, NamespacedName{Namespace: "acid", Name: "cluster"}, []byte(`"acid/cluster"`)},
|
||||
{`/name`, NamespacedName{Namespace: "default", Name: "name"}, []byte(`"default/name"`)},
|
||||
{`test`, NamespacedName{Namespace: "default", Name: "test"}, []byte(`"default/test"`)},
|
||||
{`acid/cluster`, NamespacedName{Namespace: mockOperatorNamespace, Name: "cluster"}, []byte(`"acid/cluster"`)},
|
||||
{`/name`, NamespacedName{Namespace: mockOperatorNamespace, Name: "name"}, []byte(`"acid/name"`)},
|
||||
{`test`, NamespacedName{Namespace: mockOperatorNamespace, Name: "test"}, []byte(`"acid/test"`)},
|
||||
}
|
||||
|
||||
var nnErr = []string{"test/", "/", "", "//"}
|
||||
|
||||
func TestNamespacedNameDecode(t *testing.T) {
|
||||
|
||||
for _, tt := range nnTests {
|
||||
var actual NamespacedName
|
||||
err := actual.Decode(tt.s)
|
||||
err := actual.DecodeWorker(tt.s, mockOperatorNamespace)
|
||||
if err != nil {
|
||||
t.Errorf("decode error: %v", err)
|
||||
}
|
||||
@@ -28,6 +33,7 @@ func TestNamespacedNameDecode(t *testing.T) {
|
||||
t.Errorf("expected: %v, got %#v", tt.expected, actual)
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func TestNamespacedNameMarshal(t *testing.T) {
|
||||
@@ -47,7 +53,7 @@ func TestNamespacedNameMarshal(t *testing.T) {
|
||||
func TestNamespacedNameError(t *testing.T) {
|
||||
for _, tt := range nnErr {
|
||||
var actual NamespacedName
|
||||
err := actual.Decode(tt)
|
||||
err := actual.DecodeWorker(tt, mockOperatorNamespace)
|
||||
if err == nil {
|
||||
t.Errorf("error expected for %q, got: %#v", tt, actual)
|
||||
}
|
||||
|
||||
@@ -33,6 +33,7 @@ type KubernetesClient struct {
|
||||
v1core.ConfigMapsGetter
|
||||
v1core.NodesGetter
|
||||
v1core.NamespacesGetter
|
||||
v1core.ServiceAccountsGetter
|
||||
v1beta1.StatefulSetsGetter
|
||||
policyv1beta1.PodDisruptionBudgetsGetter
|
||||
apiextbeta1.CustomResourceDefinitionsGetter
|
||||
@@ -73,6 +74,7 @@ func NewFromConfig(cfg *rest.Config) (KubernetesClient, error) {
|
||||
kubeClient.ServicesGetter = client.CoreV1()
|
||||
kubeClient.EndpointsGetter = client.CoreV1()
|
||||
kubeClient.SecretsGetter = client.CoreV1()
|
||||
kubeClient.ServiceAccountsGetter = client.CoreV1()
|
||||
kubeClient.ConfigMapsGetter = client.CoreV1()
|
||||
kubeClient.PersistentVolumeClaimsGetter = client.CoreV1()
|
||||
kubeClient.PersistentVolumesGetter = client.CoreV1()
|
||||
|
||||
@@ -5,14 +5,39 @@ import (
|
||||
"time"
|
||||
)
|
||||
|
||||
// Retry calls ConditionFunc until it returns boolean true, a timeout expires or an error occurs.
|
||||
type RetryTicker interface {
|
||||
Stop()
|
||||
Tick()
|
||||
}
|
||||
|
||||
type Ticker struct {
|
||||
ticker *time.Ticker
|
||||
}
|
||||
|
||||
func (t *Ticker) Stop() { t.ticker.Stop() }
|
||||
|
||||
func (t *Ticker) Tick() { <-t.ticker.C }
|
||||
|
||||
// Retry calls ConditionFunc until either:
|
||||
// * it returns boolean true
|
||||
// * a timeout expires
|
||||
// * an error occurs
|
||||
func Retry(interval time.Duration, timeout time.Duration, f func() (bool, error)) error {
|
||||
//TODO: make the retry exponential
|
||||
if timeout < interval {
|
||||
return fmt.Errorf("timout(%s) should be greater than interval(%v)", timeout, interval)
|
||||
}
|
||||
tick := &Ticker{time.NewTicker(interval)}
|
||||
return RetryWorker(interval, timeout, tick, f)
|
||||
}
|
||||
|
||||
func RetryWorker(
|
||||
interval time.Duration,
|
||||
timeout time.Duration,
|
||||
tick RetryTicker,
|
||||
f func() (bool, error)) error {
|
||||
|
||||
maxRetries := int(timeout / interval)
|
||||
tick := time.NewTicker(interval)
|
||||
defer tick.Stop()
|
||||
|
||||
for i := 0; ; i++ {
|
||||
@@ -26,7 +51,7 @@ func Retry(interval time.Duration, timeout time.Duration, f func() (bool, error)
|
||||
if i+1 == maxRetries {
|
||||
break
|
||||
}
|
||||
<-tick.C
|
||||
tick.Tick()
|
||||
}
|
||||
return fmt.Errorf("still failing after %d retries", maxRetries)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
package retryutil
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"testing"
|
||||
)
|
||||
|
||||
type mockTicker struct {
|
||||
test *testing.T
|
||||
counter int
|
||||
}
|
||||
|
||||
func (t *mockTicker) Stop() {}
|
||||
|
||||
func (t *mockTicker) Tick() {
|
||||
t.counter += 1
|
||||
}
|
||||
|
||||
func TestRetryWorkerSuccess(t *testing.T) {
|
||||
tick := &mockTicker{t, 0}
|
||||
result := RetryWorker(10, 20, tick, func() (bool, error) {
|
||||
return true, nil
|
||||
})
|
||||
|
||||
if result != nil {
|
||||
t.Errorf("Wrong result, expected: %#v, got: %#v", nil, result)
|
||||
}
|
||||
|
||||
if tick.counter != 0 {
|
||||
t.Errorf("Ticker was started once, but it shouldn't be")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetryWorkerOneFalse(t *testing.T) {
|
||||
var counter = 0
|
||||
|
||||
tick := &mockTicker{t, 0}
|
||||
result := RetryWorker(1, 3, tick, func() (bool, error) {
|
||||
counter += 1
|
||||
|
||||
if counter <= 1 {
|
||||
return false, nil
|
||||
} else {
|
||||
return true, nil
|
||||
}
|
||||
})
|
||||
|
||||
if result != nil {
|
||||
t.Errorf("Wrong result, expected: %#v, got: %#v", nil, result)
|
||||
}
|
||||
|
||||
if tick.counter != 1 {
|
||||
t.Errorf("Ticker was started %#v, but supposed to be just once", tick.counter)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetryWorkerError(t *testing.T) {
|
||||
fail := errors.New("Error")
|
||||
|
||||
tick := &mockTicker{t, 0}
|
||||
result := RetryWorker(1, 3, tick, func() (bool, error) {
|
||||
return false, fail
|
||||
})
|
||||
|
||||
if result != fail {
|
||||
t.Errorf("Wrong result, expected: %#v, got: %#v", fail, result)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user