Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
224 changes: 224 additions & 0 deletions operator/pkg/helper/resource_management.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ import (
"k8s.io/apimachinery/pkg/api/equality"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
Expand Down Expand Up @@ -536,6 +537,98 @@ func (rm *ResourceManager) CreatePvcIfNotExists(pvc *corev1.PersistentVolumeClai
return nil
}

func (rm *ResourceManager) CreateOrResizePvc(pvc *corev1.PersistentVolumeClaim) (bool, error) {
foundPvc := &corev1.PersistentVolumeClaim{}

err := rm.kubeClient.Get(
context.TODO(),
types.NamespacedName{
Name: pvc.Name,
Namespace: pvc.Namespace,
},
foundPvc,
)

if errors.IsNotFound(err) {
logger.Info(fmt.Sprintf("Creating %s PVC", pvc.Name))

if err := rm.kubeClient.Create(context.TODO(), pvc); err != nil {
return false, err
}

return false, nil
}

if err != nil {
return false, err
}

currentSize := foundPvc.Spec.Resources.Requests[corev1.ResourceStorage]
desiredSize := pvc.Spec.Resources.Requests[corev1.ResourceStorage]

resizeInProgress := false
changed := false

switch desiredSize.Cmp(currentSize) {
case 1:
logger.Info(fmt.Sprintf("Expanding PVC %s from %s to %s", pvc.Name, currentSize.String(), desiredSize.String()))

foundPvc.Spec.Resources.Requests[corev1.ResourceStorage] = desiredSize

resizeInProgress = true
changed = true

case 0:
capacity := foundPvc.Status.Capacity[corev1.ResourceStorage]

if !capacity.IsZero() && capacity.Cmp(desiredSize) < 0 {
logger.Info(fmt.Sprintf("PVC %s resize is still in progress: requested=%s capacity=%s", pvc.Name, desiredSize.String(), capacity.String()))

resizeInProgress = true
Comment thread
Tvion marked this conversation as resolved.
}

case -1:
return false, fmt.Errorf(
"PVC %s shrinking from %s to %s is not supported",
pvc.Name,
currentSize.String(),
desiredSize.String(),
)
}

// Preserve existing PVC behavior.
if len(foundPvc.OwnerReferences) > 0 {
foundPvc.OwnerReferences = nil
changed = true
}

// Apply desired annotations only when they differ.
if pvc.Annotations != nil {
if foundPvc.Annotations == nil {
foundPvc.Annotations = make(map[string]string)
}

for key, value := range pvc.Annotations {
if foundPvc.Annotations[key] != value {
foundPvc.Annotations[key] = value
changed = true
}
}
}

// Nothing in PVC spec/metadata changed.
// Do not perform unnecessary Kubernetes Update().
if !changed {
return resizeInProgress, nil
}

if err := rm.kubeClient.Update(context.TODO(), foundPvc); err != nil {
return false, err
}

return resizeInProgress, nil
}

func (rm *ResourceManager) CreateSecretIfNotExists(secret *corev1.Secret) error {
foundSecret := &corev1.Secret{}
err := rm.kubeClient.Get(context.TODO(), types.NamespacedName{
Expand Down Expand Up @@ -1075,3 +1168,134 @@ func (rm *ResourceManager) commonLabels(name string) map[string]string {

return labels
}

func (rm *ResourceManager) WaitForPvcResizeState(pvcName string, namespace string, desiredSize resource.Quantity) (bool, error) {
restartRequired := false

err := wait.PollUntilContextTimeout(context.Background(), time.Second, 2*time.Minute, true,
func(ctx context.Context) (bool, error) {
currentPvc := &corev1.PersistentVolumeClaim{}

if err := rm.kubeClient.Get(
ctx,
types.NamespacedName{
Name: pvcName,
Namespace: namespace,
},
currentPvc,
); err != nil {
return false, err
}

capacity := currentPvc.Status.Capacity[corev1.ResourceStorage]

if capacity.Cmp(desiredSize) >= 0 {
return true, nil
}

for _, condition := range currentPvc.Status.Conditions {
if condition.Type == corev1.PersistentVolumeClaimFileSystemResizePending &&
condition.Status == corev1.ConditionTrue {

restartRequired = true
return true, nil
}
}

return false, nil
},
)

return restartRequired, err
}

func (rm *ResourceManager) WaitForPvcCapacity(pvcName string, namespace string, desiredSize resource.Quantity, timeout time.Duration) (bool, error) {
resized := false

err := wait.PollUntilContextTimeout(context.Background(), time.Second, timeout, true,
func(ctx context.Context) (bool, error) {
pvc := &corev1.PersistentVolumeClaim{}

if err := rm.kubeClient.Get(
ctx,
types.NamespacedName{
Name: pvcName,
Namespace: namespace,
},
pvc,
); err != nil {
logger.Warn(
fmt.Sprintf("Failed to get PVC %s while waiting for resize, retrying", pvcName),
zap.Error(err),
)

return false, nil
}

capacity := pvc.Status.Capacity[corev1.ResourceStorage]

if capacity.Cmp(desiredSize) >= 0 {
logger.Info(fmt.Sprintf("PVC %s resize completed: capacity=%s", pvcName, capacity.String()))

resized = true
return true, nil
}

logger.Info(fmt.Sprintf("Waiting for PVC %s capacity: current=%s desired=%s", pvcName, capacity.String(), desiredSize.String()))

return false, nil
},
)

if err == context.DeadlineExceeded {
return false, nil
}

if err != nil {
return false, err
}

return resized, nil
}

func (rm *ResourceManager) ScaleStatefulSet(name string, replicas int32) error {
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
sts := &appsv1.StatefulSet{}

if err := rm.kubeClient.Get(
context.TODO(),
types.NamespacedName{
Name: name,
Namespace: namespace,
},
sts,
); err != nil {
return err
}

sts.Spec.Replicas = &replicas

return rm.kubeClient.Update(context.TODO(), sts)
})
}

func (rm *ResourceManager) ScaleDeployment(name string, replicas int32) error {
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
deployment := &appsv1.Deployment{}

if err := rm.kubeClient.Get(
context.TODO(),
types.NamespacedName{
Name: name,
Namespace: namespace,
},
deployment,
); err != nil {
return err
}

deployment.Spec.Replicas = &replicas

return rm.kubeClient.Update(context.TODO(), deployment)
})
}
27 changes: 27 additions & 0 deletions operator/pkg/patroni/patroni.go
Original file line number Diff line number Diff line change
Expand Up @@ -623,6 +623,33 @@ func GenerateLDAPConfig(cr *patroniv1.PatroniCore) []string {
}
}

func Switchover(patroniURL, leader, candidate string) error {
body := map[string]string{"leader": leader, "candidate": candidate}

data, err := json.Marshal(body)
if err != nil {
return err
}

resp, err := patroniPost(http.DefaultClient, patroniURL+"switchover", bytes.NewBuffer(data))
if err != nil {
return err
}
defer resp.Body.Close()

if resp.StatusCode < http.StatusOK ||
resp.StatusCode >= http.StatusMultipleChoices {

responseBody, _ := io.ReadAll(resp.Body)

return fmt.Errorf("Patroni switchover failed: %s, response: %s", resp.Status, string(responseBody))
}

logger.Info(fmt.Sprintf("Patroni switchover from %s to %s requested successfully", leader, candidate))

return nil
}

func ReloadPatroniConfig(patroniUrl string) error {
hosts, err := getPatroniHosts(patroniUrl)
if err != nil {
Expand Down
Loading
Loading