diff --git a/apis/workloads/v1/instance_types.go b/apis/workloads/v1/instance_types.go
index 2009a52cd7e..b90f00664d1 100644
--- a/apis/workloads/v1/instance_types.go
+++ b/apis/workloads/v1/instance_types.go
@@ -153,7 +153,7 @@ type InstanceStatus2 struct {
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
// Represents the latest available observations of an instance's current state.
- // Known .status.conditions.type are: "InstanceFailure", "InstanceReady", "InstanceAvailable"
+ // Known .status.conditions.type are: "InstanceFailure", "InstanceReady", "InstanceAvailable", "Restore"
//
// +optional
Conditions []metav1.Condition `json:"conditions,omitempty"`
diff --git a/apis/workloads/v1/instanceset_types.go b/apis/workloads/v1/instanceset_types.go
index 1b7ead7a19f..f6081257067 100644
--- a/apis/workloads/v1/instanceset_types.go
+++ b/apis/workloads/v1/instanceset_types.go
@@ -738,7 +738,7 @@ const (
// InstanceFailure is added in an instance set when at least one of its instances(pods) is in a `Failed` phase.
InstanceFailure ConditionType = "InstanceFailure"
- // InstanceRestore indicates whether the initial data restore for this InstanceSet has completed.
+ // InstanceRestore indicates whether the initial data restore for this Instance or InstanceSet has completed.
InstanceRestore ConditionType = "Restore"
// InstanceUpdateRestricted represents a ConditionType that indicates updates to an InstanceSet are blocked(when the
diff --git a/config/crd/bases/workloads.kubeblocks.io_instances.yaml b/config/crd/bases/workloads.kubeblocks.io_instances.yaml
index 43bdfc62774..6df616942d8 100644
--- a/config/crd/bases/workloads.kubeblocks.io_instances.yaml
+++ b/config/crd/bases/workloads.kubeblocks.io_instances.yaml
@@ -10267,7 +10267,7 @@ spec:
conditions:
description: |-
Represents the latest available observations of an instance's current state.
- Known .status.conditions.type are: "InstanceFailure", "InstanceReady", "InstanceAvailable"
+ Known .status.conditions.type are: "InstanceFailure", "InstanceReady", "InstanceAvailable", "Restore"
items:
description: Condition contains details for one aspect of the current
state of this API Resource.
diff --git a/deploy/helm/crds/workloads.kubeblocks.io_instances.yaml b/deploy/helm/crds/workloads.kubeblocks.io_instances.yaml
index 43bdfc62774..6df616942d8 100644
--- a/deploy/helm/crds/workloads.kubeblocks.io_instances.yaml
+++ b/deploy/helm/crds/workloads.kubeblocks.io_instances.yaml
@@ -10267,7 +10267,7 @@ spec:
conditions:
description: |-
Represents the latest available observations of an instance's current state.
- Known .status.conditions.type are: "InstanceFailure", "InstanceReady", "InstanceAvailable"
+ Known .status.conditions.type are: "InstanceFailure", "InstanceReady", "InstanceAvailable", "Restore"
items:
description: Condition contains details for one aspect of the current
state of this API Resource.
diff --git a/docs/developer_docs/api-reference/cluster.md b/docs/developer_docs/api-reference/cluster.md
index 7f67f1444aa..0b797c27069 100644
--- a/docs/developer_docs/api-reference/cluster.md
+++ b/docs/developer_docs/api-reference/cluster.md
@@ -19327,7 +19327,7 @@ ConditionStatus will be True if all its instances(pods) are in a Ready condition
Or, a NotReady reason with not ready instances encoded in the Message filed will be set.
"Restore" |
-InstanceRestore indicates whether the initial data restore for this InstanceSet has completed.
+ | InstanceRestore indicates whether the initial data restore for this Instance or InstanceSet has completed.
|
"InstanceUpdateRestricted" |
InstanceUpdateRestricted represents a ConditionType that indicates updates to an InstanceSet are blocked(when the
@@ -20728,7 +20728,7 @@ Instance’s generation, which is updated on mutation by the API Server.
|
(Optional)
Represents the latest available observations of an instance’s current state.
-Known .status.conditions.type are: “InstanceFailure”, “InstanceReady”, “InstanceAvailable”
+Known .status.conditions.type are: “InstanceFailure”, “InstanceReady”, “InstanceAvailable”, “Restore”
|
diff --git a/pkg/controller/instance/reconciler_status.go b/pkg/controller/instance/reconciler_status.go
index 79400933b08..cd898678eb2 100644
--- a/pkg/controller/instance/reconciler_status.go
+++ b/pkg/controller/instance/reconciler_status.go
@@ -20,6 +20,9 @@ along with this program. If not, see .
package instance
import (
+ "fmt"
+ "sort"
+ "strings"
"time"
corev1 "k8s.io/api/core/v1"
@@ -50,6 +53,7 @@ func (r *statusReconciler) PreCondition(tree *kubebuilderx.ObjectTree) *kubebuil
func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilderx.Result, error) {
inst := tree.GetRoot().(*workloads.Instance)
+ r.reconcileRestoreCondition(tree, inst)
obj, err := tree.Get(podObj(inst))
if err != nil {
@@ -122,6 +126,87 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
return kubebuilderx.Continue, nil
}
+func (r *statusReconciler) reconcileRestoreCondition(tree *kubebuilderx.ObjectTree, inst *workloads.Instance) {
+ restoreCond := meta.FindStatusCondition(inst.Status.Conditions, string(workloads.InstanceRestore))
+ if restoreCond != nil && (restoreCond.Status == metav1.ConditionTrue || restoreCond.Status == metav1.ConditionFalse) {
+ return
+ }
+ condition := r.buildRestoreCondition(tree, inst)
+ if condition == nil {
+ meta.RemoveStatusCondition(&inst.Status.Conditions, string(workloads.InstanceRestore))
+ return
+ }
+ meta.SetStatusCondition(&inst.Status.Conditions, *condition)
+}
+
+func (r *statusReconciler) buildRestoreCondition(tree *kubebuilderx.ObjectTree, inst *workloads.Instance) *metav1.Condition {
+ expectedPVCNames := make(map[string]struct{})
+ for i := range inst.Spec.VolumeClaimTemplates {
+ vct := &inst.Spec.VolumeClaimTemplates[i]
+ if vct.Annotations[constant.RestoreSourceKindAnnotationKey] == "" {
+ continue
+ }
+ expectedPVCNames[intctrlutil.ComposePVCName(corev1.PersistentVolumeClaim{ObjectMeta: vct.ObjectMeta}, inst.Spec.InstanceSetName, inst.Name)] = struct{}{}
+ }
+ if len(expectedPVCNames) == 0 {
+ return nil
+ }
+
+ pvcsByName := r.persistentVolumeClaimsByName(tree)
+ completed := 0
+ waiting := make([]string, 0, len(expectedPVCNames))
+ for name := range expectedPVCNames {
+ pvc := pvcsByName[name]
+ if pvc == nil {
+ waiting = append(waiting, name)
+ continue
+ }
+ cond := findPVCRestoreCondition(pvc)
+ if cond == nil || cond.Status == corev1.ConditionUnknown {
+ waiting = append(waiting, name)
+ continue
+ }
+ if cond.Status == corev1.ConditionFalse {
+ return &metav1.Condition{
+ Type: string(workloads.InstanceRestore),
+ Status: metav1.ConditionFalse,
+ ObservedGeneration: inst.Generation,
+ Reason: workloads.ReasonRestoreFailed,
+ Message: fmt.Sprintf("PVC %s restore failed: %s", pvc.Name, cond.Message),
+ }
+ }
+ if cond.Status == corev1.ConditionTrue {
+ completed++
+ }
+ }
+ if completed == len(expectedPVCNames) {
+ return &metav1.Condition{
+ Type: string(workloads.InstanceRestore),
+ Status: metav1.ConditionTrue,
+ ObservedGeneration: inst.Generation,
+ Reason: workloads.ReasonRestoreCompleted,
+ Message: "All initial restore PVCs have completed",
+ }
+ }
+ sort.Strings(waiting)
+ return &metav1.Condition{
+ Type: string(workloads.InstanceRestore),
+ Status: metav1.ConditionUnknown,
+ ObservedGeneration: inst.Generation,
+ Reason: workloads.ReasonRestoreRunning,
+ Message: fmt.Sprintf("Waiting for initial restore PVCs to complete: %s", strings.Join(waiting, ",")),
+ }
+}
+
+func findPVCRestoreCondition(pvc *corev1.PersistentVolumeClaim) *corev1.PersistentVolumeClaimCondition {
+ for i := range pvc.Status.Conditions {
+ if string(pvc.Status.Conditions[i].Type) == string(workloads.InstanceRestore) {
+ return &pvc.Status.Conditions[i]
+ }
+ }
+ return nil
+}
+
// An absent or terminating Pod cannot provide a valid runtime observation. Clear every Pod-derived field and
// condition together so values from the previous Pod do not survive the lifecycle transition.
func (r *statusReconciler) setPodUnavailableStatus(inst *workloads.Instance, state workloads.InstanceCurrentState, name, revision string) {
diff --git a/pkg/controller/instance/reconciler_status_test.go b/pkg/controller/instance/reconciler_status_test.go
index 4727423959d..50092e0f6b4 100644
--- a/pkg/controller/instance/reconciler_status_test.go
+++ b/pkg/controller/instance/reconciler_status_test.go
@@ -134,6 +134,76 @@ func TestStatusReconcilerPublishesInstanceCurrentState(t *testing.T) {
}
}
+func TestStatusReconcilerAggregatesRestorePVCConditionsWithoutPod(t *testing.T) {
+ newFixture := func() (*workloads.Instance, *kubebuilderx.ObjectTree, string) {
+ claim := corev1.PersistentVolumeClaimTemplate{ObjectMeta: metav1.ObjectMeta{
+ Name: "data",
+ Annotations: map[string]string{
+ constant.RestoreSourceKindAnnotationKey: "Backup",
+ },
+ }}
+ inst := &workloads.Instance{
+ ObjectMeta: metav1.ObjectMeta{Name: "demo-0", Namespace: "default", Generation: 2},
+ Spec: workloads.InstanceSpec{
+ InstanceSetName: "demo",
+ VolumeClaimTemplates: []corev1.PersistentVolumeClaimTemplate{claim},
+ },
+ }
+ tree := kubebuilderx.NewObjectTree()
+ tree.SetRoot(inst)
+ return inst, tree, intctrlutil.ComposePVCName(corev1.PersistentVolumeClaim{ObjectMeta: claim.ObjectMeta}, "demo", "demo-0")
+ }
+ addPVC := func(t *testing.T, tree *kubebuilderx.ObjectTree, name string, status corev1.ConditionStatus) {
+ t.Helper()
+ pvc := &corev1.PersistentVolumeClaim{
+ ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "default"},
+ Status: corev1.PersistentVolumeClaimStatus{Conditions: []corev1.PersistentVolumeClaimCondition{{
+ Type: corev1.PersistentVolumeClaimConditionType(workloads.InstanceRestore),
+ Status: status,
+ Message: "restore result",
+ }}},
+ }
+ if err := tree.Add(pvc); err != nil {
+ t.Fatal(err)
+ }
+ }
+
+ t.Run("waits for the PVC before the Pod exists", func(t *testing.T) {
+ inst, tree, _ := newFixture()
+ if _, err := NewStatusReconciler().Reconcile(tree); err != nil {
+ t.Fatal(err)
+ }
+ cond := meta.FindStatusCondition(inst.Status.Conditions, string(workloads.InstanceRestore))
+ if cond == nil || cond.Status != metav1.ConditionUnknown || cond.Reason != workloads.ReasonRestoreRunning {
+ t.Fatalf("unexpected Restore condition: %#v", cond)
+ }
+ })
+
+ t.Run("publishes completed", func(t *testing.T) {
+ inst, tree, pvcName := newFixture()
+ addPVC(t, tree, pvcName, corev1.ConditionTrue)
+ if _, err := NewStatusReconciler().Reconcile(tree); err != nil {
+ t.Fatal(err)
+ }
+ cond := meta.FindStatusCondition(inst.Status.Conditions, string(workloads.InstanceRestore))
+ if cond == nil || cond.Status != metav1.ConditionTrue || cond.Reason != workloads.ReasonRestoreCompleted {
+ t.Fatalf("unexpected Restore condition: %#v", cond)
+ }
+ })
+
+ t.Run("publishes failure first", func(t *testing.T) {
+ inst, tree, pvcName := newFixture()
+ addPVC(t, tree, pvcName, corev1.ConditionFalse)
+ if _, err := NewStatusReconciler().Reconcile(tree); err != nil {
+ t.Fatal(err)
+ }
+ cond := meta.FindStatusCondition(inst.Status.Conditions, string(workloads.InstanceRestore))
+ if cond == nil || cond.Status != metav1.ConditionFalse || cond.Reason != workloads.ReasonRestoreFailed {
+ t.Fatalf("unexpected Restore condition: %#v", cond)
+ }
+ })
+}
+
func TestStatusReconcilerKeepsUpToDateFalseUntilPVCExpansionCompletes(t *testing.T) {
claim := corev1.PersistentVolumeClaimTemplate{
ObjectMeta: metav1.ObjectMeta{Name: "data"},
diff --git a/pkg/controller/instanceset2/reconciler_status.go b/pkg/controller/instanceset2/reconciler_status.go
index f24724908b5..5a12779a48b 100644
--- a/pkg/controller/instanceset2/reconciler_status.go
+++ b/pkg/controller/instanceset2/reconciler_status.go
@@ -175,6 +175,10 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
meta.RemoveStatusCondition(&its.Status.Conditions, string(workloads.InstanceFailure))
}
+ if err = r.reconcileRestoreCondition(tree, its, instanceList); err != nil {
+ return kubebuilderx.Continue, err
+ }
+
// 4. set instance status
if err := setInstanceStatus(tree, its, instanceList); err != nil {
return kubebuilderx.Continue, err
@@ -186,6 +190,93 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
return kubebuilderx.Continue, nil
}
+func (r *statusReconciler) reconcileRestoreCondition(tree *kubebuilderx.ObjectTree, its *workloads.InstanceSet, instances []*workloads.Instance) error {
+ restoreCond := meta.FindStatusCondition(its.Status.Conditions, string(workloads.InstanceRestore))
+ if restoreCond != nil && (restoreCond.Status == metav1.ConditionTrue || restoreCond.Status == metav1.ConditionFalse) {
+ return nil
+ }
+ condition, err := buildRestoreCondition(tree, its, instances)
+ if err != nil {
+ return err
+ }
+ if condition == nil {
+ meta.RemoveStatusCondition(&its.Status.Conditions, string(workloads.InstanceRestore))
+ return nil
+ }
+ meta.SetStatusCondition(&its.Status.Conditions, *condition)
+ return nil
+}
+
+func buildRestoreCondition(tree *kubebuilderx.ObjectTree, its *workloads.InstanceSet, instances []*workloads.Instance) (*metav1.Condition, error) {
+ desiredInstances, _, err := buildDesiredInstancesByName(tree, its)
+ if err != nil {
+ return nil, err
+ }
+ expectedNames := sets.New[string]()
+ for name, inst := range desiredInstances {
+ for i := range inst.Spec.VolumeClaimTemplates {
+ if inst.Spec.VolumeClaimTemplates[i].Annotations[constant.RestoreSourceKindAnnotationKey] != "" {
+ expectedNames.Insert(name)
+ break
+ }
+ }
+ }
+ if expectedNames.Len() == 0 {
+ return nil, nil
+ }
+
+ instancesByName := make(map[string]*workloads.Instance, len(instances))
+ for _, inst := range instances {
+ instancesByName[inst.Name] = inst
+ }
+ completed := 0
+ waiting := sets.New[string]()
+ for name := range expectedNames {
+ inst := instancesByName[name]
+ if inst == nil {
+ waiting.Insert(name)
+ continue
+ }
+ cond := meta.FindStatusCondition(inst.Status.Conditions, string(workloads.InstanceRestore))
+ if cond == nil || cond.Status == metav1.ConditionUnknown {
+ waiting.Insert(name)
+ continue
+ }
+ if cond.Status == metav1.ConditionFalse {
+ return &metav1.Condition{
+ Type: string(workloads.InstanceRestore),
+ Status: metav1.ConditionFalse,
+ ObservedGeneration: its.Generation,
+ Reason: workloads.ReasonRestoreFailed,
+ Message: fmt.Sprintf("Instance %s restore failed: %s", inst.Name, cond.Message),
+ }, nil
+ }
+ if cond.Status == metav1.ConditionTrue {
+ completed++
+ }
+ }
+ if completed == expectedNames.Len() {
+ return &metav1.Condition{
+ Type: string(workloads.InstanceRestore),
+ Status: metav1.ConditionTrue,
+ ObservedGeneration: its.Generation,
+ Reason: workloads.ReasonRestoreCompleted,
+ Message: "All initial restore Instances have completed",
+ }, nil
+ }
+ message, err := buildConditionMessageWithNames(waiting.UnsortedList())
+ if err != nil {
+ return nil, err
+ }
+ return &metav1.Condition{
+ Type: string(workloads.InstanceRestore),
+ Status: metav1.ConditionUnknown,
+ ObservedGeneration: its.Generation,
+ Reason: workloads.ReasonRestoreRunning,
+ Message: fmt.Sprintf("Waiting for initial restore Instances to complete: %s", message),
+ }, nil
+}
+
func getInstanceTemplateName(inst *workloads.Instance) string {
if inst.Labels == nil {
return ""
diff --git a/pkg/controller/instanceset2/reconciler_status_test.go b/pkg/controller/instanceset2/reconciler_status_test.go
index 576dd258a94..d2897809006 100644
--- a/pkg/controller/instanceset2/reconciler_status_test.go
+++ b/pkg/controller/instanceset2/reconciler_status_test.go
@@ -24,6 +24,7 @@ import (
"testing"
corev1 "k8s.io/api/core/v1"
+ "k8s.io/apimachinery/pkg/api/meta"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/utils/ptr"
@@ -35,6 +36,86 @@ import (
"github.com/apecloud/kubeblocks/pkg/controller/revisionmap"
)
+func TestStatusReconcilerAggregatesInstanceRestoreConditions(t *testing.T) {
+ newFixture := func() (*workloads.InstanceSet, *kubebuilderx.ObjectTree) {
+ its := &workloads.InstanceSet{
+ ObjectMeta: metav1.ObjectMeta{Name: "demo", Namespace: "default", Generation: 3},
+ Spec: workloads.InstanceSetSpec{
+ Replicas: ptr.To[int32](1),
+ Selector: &metav1.LabelSelector{MatchLabels: map[string]string{"app": "demo"}},
+ VolumeClaimTemplates: []corev1.PersistentVolumeClaim{{ObjectMeta: metav1.ObjectMeta{
+ Name: "data",
+ Annotations: map[string]string{
+ constant.RestoreSourceKindAnnotationKey: "Backup",
+ },
+ }}},
+ },
+ }
+ tree := kubebuilderx.NewObjectTree()
+ tree.SetRoot(its)
+ return its, tree
+ }
+ newInstance := func(status metav1.ConditionStatus) *workloads.Instance {
+ return &workloads.Instance{
+ ObjectMeta: metav1.ObjectMeta{Name: "demo-0", Namespace: "default"},
+ Status: workloads.InstanceStatus2{Conditions: []metav1.Condition{{
+ Type: string(workloads.InstanceRestore),
+ Status: status,
+ Message: "restore result",
+ }}},
+ }
+ }
+
+ t.Run("waits for the desired Instance", func(t *testing.T) {
+ its, tree := newFixture()
+ cond, err := buildRestoreCondition(tree, its, nil)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if cond == nil || cond.Status != metav1.ConditionUnknown || cond.Reason != workloads.ReasonRestoreRunning {
+ t.Fatalf("unexpected Restore condition: %#v", cond)
+ }
+ })
+
+ t.Run("publishes completed", func(t *testing.T) {
+ its, tree := newFixture()
+ cond, err := buildRestoreCondition(tree, its, []*workloads.Instance{newInstance(metav1.ConditionTrue)})
+ if err != nil {
+ t.Fatal(err)
+ }
+ if cond == nil || cond.Status != metav1.ConditionTrue || cond.Reason != workloads.ReasonRestoreCompleted {
+ t.Fatalf("unexpected Restore condition: %#v", cond)
+ }
+ })
+
+ t.Run("publishes failure first", func(t *testing.T) {
+ its, tree := newFixture()
+ cond, err := buildRestoreCondition(tree, its, []*workloads.Instance{newInstance(metav1.ConditionFalse)})
+ if err != nil {
+ t.Fatal(err)
+ }
+ if cond == nil || cond.Status != metav1.ConditionFalse || cond.Reason != workloads.ReasonRestoreFailed {
+ t.Fatalf("unexpected Restore condition: %#v", cond)
+ }
+ })
+
+ t.Run("keeps a terminal failure", func(t *testing.T) {
+ its, tree := newFixture()
+ meta.SetStatusCondition(&its.Status.Conditions, metav1.Condition{
+ Type: string(workloads.InstanceRestore),
+ Status: metav1.ConditionFalse,
+ Reason: workloads.ReasonRestoreFailed,
+ })
+ if err := (&statusReconciler{}).reconcileRestoreCondition(tree, its, []*workloads.Instance{newInstance(metav1.ConditionTrue)}); err != nil {
+ t.Fatal(err)
+ }
+ cond := meta.FindStatusCondition(its.Status.Conditions, string(workloads.InstanceRestore))
+ if cond == nil || cond.Status != metav1.ConditionFalse {
+ t.Fatalf("terminal Restore condition was overwritten: %#v", cond)
+ }
+ })
+}
+
func TestSetInstanceStatusReadsCurrentStateFromInstance(t *testing.T) {
its := &workloads.InstanceSet{
ObjectMeta: metav1.ObjectMeta{Name: "demo", Namespace: "default", Generation: 3},