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
66 changes: 66 additions & 0 deletions pkg/tnf/hack/test-transition-events.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
#!/usr/bin/bash
# Test script to validate transition event format on a live cluster.
# Creates mock events identical to what our code produces, then verifies
# they appear correctly in `oc get events`.
#
# Usage: ./hack/test-transition-events.sh
# Requires: oc/kubectl with cluster-admin access

set -euo pipefail

NAMESPACE="openshift-etcd"
TIMESTAMP=$(date -u +"%Y-%m-%dT%H:%M:%SZ")

events=(
"EtcdTransitionAuthCompleted|PCS authentication completed on all nodes"
"EtcdTransitionClusterConfigured|Pacemaker cluster configured successfully"
"EtcdTransitionFencingConfigured|STONITH fencing configured successfully"
"EtcdTransitionEtcdResourceCreated|Pacemaker etcd resource agent (podman-etcd) configured"
"EtcdTransitionConstraintsConfigured|Pacemaker ordering and colocation constraints configured"
"EtcdTransitionStarted|Etcd transition from CEO-controlled to pacemaker-controlled has started"
"EtcdTransitionWaitingForRemoval|Waiting for CEO to remove static etcd container from all nodes"
"EtcdTransitionStaticContainerRemoved|Static etcd container removed from all nodes, revision is stable"
"EtcdTransitionCompleted|Etcd transition to pacemaker-controlled has completed"
)

echo "Creating ${#events[@]} test transition events in ${NAMESPACE}..."

for entry in "${events[@]}"; do
reason="${entry%%|*}"
message="${entry#*|}"
name="test-${reason,,}-$(date +%s%N)"

oc apply -f - <<EOF
apiVersion: v1
kind: Event
metadata:
name: ${name}
namespace: ${NAMESPACE}
labels:
tnf-test: "true"
involvedObject:
kind: Job
name: tnf-setup-job
namespace: ${NAMESPACE}
apiVersion: batch/v1
reason: ${reason}
message: "${message}"
type: Normal
source:
component: tnf-setup-runner
firstTimestamp: "${TIMESTAMP}"
lastTimestamp: "${TIMESTAMP}"
count: 1
EOF

echo " Created: ${reason}"
done

echo ""
echo "Verifying events..."
echo ""
oc get events -n "${NAMESPACE}" --field-selector reason!=="" --sort-by='.lastTimestamp' | grep -i "EtcdTransition" || echo "No EtcdTransition events found!"

echo ""
echo "Cleanup: oc delete events -n ${NAMESPACE} -l tnf-test=true"
echo " (events auto-expire after ~1 hour)"
26 changes: 22 additions & 4 deletions pkg/tnf/pkg/etcd/etcd.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,15 @@ import (
operatorv1 "github.com/openshift/api/operator/v1"
"github.com/openshift/library-go/pkg/operator/v1helpers"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go/kubernetes"
"k8s.io/klog/v2"

"github.com/openshift/cluster-etcd-operator/pkg/operator/ceohelpers"
"github.com/openshift/cluster-etcd-operator/pkg/tnf/pkg/tools"
)

// RemoveStaticContainer informs CEO to remove its etcd container
func RemoveStaticContainer(ctx context.Context, operatorClient v1helpers.StaticPodOperatorClient) error {
func RemoveStaticContainer(ctx context.Context, operatorClient v1helpers.StaticPodOperatorClient, kubeClient kubernetes.Interface) error {
klog.Info("Signaling CEO that TNF setup is ready for etcd container transition")

_, _, err := v1helpers.UpdateStatus(ctx, operatorClient, v1helpers.UpdateConditionFn(operatorv1.OperatorCondition{
Expand All @@ -28,8 +30,11 @@ func RemoveStaticContainer(ctx context.Context, operatorClient v1helpers.StaticP
return fmt.Errorf("error while updating ExternalEtcdReadyForTransition operator condition: %w", err)
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionStarted",
"Etcd transition from CEO-controlled to pacemaker-controlled has started")

// Wait for CEO to respond by removing static containers
err = waitForStaticContainerRemoved(ctx, operatorClient)
err = waitForStaticContainerRemoved(ctx, operatorClient, kubeClient)
if err != nil {
klog.Error(err, "Failed to wait for etcd container transition")
return err
Expand All @@ -39,15 +44,21 @@ func RemoveStaticContainer(ctx context.Context, operatorClient v1helpers.StaticP
}

// waitForStaticContainerRemoved waits until the static etcd container has been removed
func waitForStaticContainerRemoved(ctx context.Context, operatorClient v1helpers.StaticPodOperatorClient) error {
func waitForStaticContainerRemoved(ctx context.Context, operatorClient v1helpers.StaticPodOperatorClient, kubeClient kubernetes.Interface) error {
klog.Info("Wait for static etcd removed")

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionWaitingForRemoval",
"Waiting for CEO to remove static etcd container from all nodes")

// the container is removed when all nodes run the latest revision
err := WaitForStableRevision(ctx, operatorClient)
if err != nil {
return err
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionStaticContainerRemoved",
"Static etcd container removed from all nodes, revision is stable")

// Update the operator status to indicate that the transition has completed.
// As soon as the etcd container is removed, this operator won't be able to update this status
// unless the etcd container is restarted by the pacemaker resource agent.
Expand All @@ -58,7 +69,14 @@ func waitForStaticContainerRemoved(ctx context.Context, operatorClient v1helpers
Message: "pacemaker's resource agent is now running the etcd container",
}))

return err
if err != nil {
return err
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionCompleted",
"Etcd transition to pacemaker-controlled has completed, pacemaker's resource agent is now running the etcd container")

return nil
}

// WaitForStableRevision waits until all nodes run the latest available revision
Expand Down
74 changes: 31 additions & 43 deletions pkg/tnf/pkg/etcd/etcd_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,54 +8,47 @@ import (
"github.com/openshift/library-go/pkg/operator/v1helpers"
"github.com/stretchr/testify/require"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes/fake"

"github.com/openshift/cluster-etcd-operator/pkg/operator/ceohelpers"
)

type args struct {
ctx context.Context
operatorClient v1helpers.StaticPodOperatorClient
}

func TestRemoveStaticContainer(t *testing.T) {
tests := []struct {
name string
args args
wantErr bool
}{
{
name: "sets ExternalEtcdReadyForTransition condition",
args: getArgs(),
wantErr: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if err := RemoveStaticContainer(tt.args.ctx, tt.args.operatorClient); (err != nil) != tt.wantErr {
t.Errorf("RemoveStaticContainer() error = %v, wantErr %v", err, tt.wantErr)
}
ctx := context.Background()
operatorClient := newFakeOperatorClient()
kubeClient := fake.NewSimpleClientset()

// Verify that the operator condition was set
_, status, _, err := tt.args.operatorClient.GetStaticPodOperatorState()
require.NoError(t, err, "Failed to get static pod operator state")
isSet := v1helpers.IsOperatorConditionTrue(status.Conditions, ceohelpers.OperatorConditionExternalEtcdReadyForTransition)
require.True(t, isSet, "Expected ReadyForEtcdContainerRemoval condition to be set to True")
})
}
}
err := RemoveStaticContainer(ctx, operatorClient, kubeClient)
require.NoError(t, err)

func getArgs() args {
ctx := context.Background()
_, status, _, err := operatorClient.GetStaticPodOperatorState()
require.NoError(t, err)
require.True(t, v1helpers.IsOperatorConditionTrue(status.Conditions, ceohelpers.OperatorConditionExternalEtcdReadyForTransition))
require.True(t, v1helpers.IsOperatorConditionTrue(status.Conditions, ceohelpers.OperatorConditionExternalEtcdHasCompletedTransition))

etcd := &v1.Etcd{
ObjectMeta: metav1.ObjectMeta{
Name: "cluster",
},
events, err := kubeClient.CoreV1().Events("openshift-etcd").List(ctx, metav1.ListOptions{})
require.NoError(t, err)

expectedReasons := map[string]bool{
"EtcdTransitionStarted": true,
"EtcdTransitionWaitingForRemoval": true,
"EtcdTransitionStaticContainerRemoved": true,
"EtcdTransitionCompleted": true,
}

// Create a fake operator client with node statuses that indicate etcd container has been removed
fakeOperatorClient := v1helpers.NewFakeStaticPodOperatorClient(
&etcd.Spec.StaticPodOperatorSpec,
require.Len(t, events.Items, len(expectedReasons))
for _, event := range events.Items {
require.True(t, expectedReasons[event.Reason], "unexpected event reason: %s", event.Reason)
require.Equal(t, "Normal", event.Type)
require.Equal(t, "tnf-setup-runner", event.Source.Component)
delete(expectedReasons, event.Reason)
}
require.Empty(t, expectedReasons, "missing events: %v", expectedReasons)
}

func newFakeOperatorClient() v1helpers.StaticPodOperatorClient {
return v1helpers.NewFakeStaticPodOperatorClient(
&v1.StaticPodOperatorSpec{},
&v1.StaticPodOperatorStatus{
OperatorStatus: v1.OperatorStatus{
LatestAvailableRevision: 1,
Expand All @@ -70,9 +63,4 @@ func getArgs() args {
nil,
nil,
)

return args{
ctx: ctx,
operatorClient: fakeOperatorClient,
}
}
50 changes: 50 additions & 0 deletions pkg/tnf/pkg/tools/events.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
package tools

import (
"context"
"fmt"
"strings"
"time"

corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/klog/v2"
)

const (
setupEventNamespace = "openshift-etcd"
setupJobName = "tnf-setup-job"
setupSourceComponent = "tnf-setup-runner"
)

// RecordSetupEvent creates a Normal Kubernetes event in openshift-etcd for a
// TNF setup lifecycle milestone. Failures are logged but never returned —
// event recording must not block the setup.
func RecordSetupEvent(ctx context.Context, kubeClient kubernetes.Interface, reason, message string) {
event := &corev1.Event{
ObjectMeta: metav1.ObjectMeta{
Name: fmt.Sprintf("tnf-%s-%d", strings.ToLower(reason), time.Now().UnixNano()),
Namespace: setupEventNamespace,
},
Comment thread
coderabbitai[bot] marked this conversation as resolved.
InvolvedObject: corev1.ObjectReference{
Kind: "Job",
Name: setupJobName,
Namespace: setupEventNamespace,
APIVersion: "batch/v1",
},
Reason: reason,
Message: message,
Type: corev1.EventTypeNormal,
Source: corev1.EventSource{Component: setupSourceComponent},
FirstTimestamp: metav1.Now(),
LastTimestamp: metav1.Now(),
Count: 1,
}

if _, err := kubeClient.CoreV1().Events(setupEventNamespace).Create(ctx, event, metav1.CreateOptions{}); err != nil {
klog.Warningf("Failed to record %s event: %v", reason, err)
} else {
klog.Infof("Recorded event: %s - %s", reason, message)
}
}
17 changes: 16 additions & 1 deletion pkg/tnf/setup/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,9 @@ func RunTnfSetup() error {
return err
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionAuthCompleted",
"PCS authentication completed on all nodes")

klog.Info("Running TNF setup")

// create tnf cluster config
Expand All @@ -100,12 +103,18 @@ func RunTnfSetup() error {
time.Sleep(5 * time.Second)
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionClusterConfigured",
"Pacemaker cluster configured successfully")

// configure stonith
err = pcs.ConfigureFencing(ctx, kubeClient, []string{cfg.NodeName1, cfg.NodeName2})
if err != nil {
return err
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionFencingConfigured",
"STONITH fencing configured successfully")

// Register pacemaker alert agents for fencing taint/untaint.
// Tolerate the scripts not being present yet, they are delivered by
// MCO which may not have rolled out at this point
Expand All @@ -125,6 +134,9 @@ func RunTnfSetup() error {
return err
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionEtcdResourceCreated",
"Pacemaker etcd resource agent (podman-etcd) configured")

// configure etcd constraints
configured, err = pcs.ConfigureConstraints(ctx)
if err != nil {
Expand All @@ -134,8 +146,11 @@ func RunTnfSetup() error {
time.Sleep(5 * time.Second)
}

tools.RecordSetupEvent(ctx, kubeClient, "EtcdTransitionConstraintsConfigured",
"Pacemaker ordering and colocation constraints configured")

// Signal CEO that TNF setup is ready for etcd container removal
err = etcd.RemoveStaticContainer(ctx, operatorClient)
err = etcd.RemoveStaticContainer(ctx, operatorClient, kubeClient)
if err != nil {
return err
}
Expand Down