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
77 changes: 77 additions & 0 deletions pkg/reconciler/common/transformers.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,10 @@ const (
runAsNonRootValue = true
allowPrivilegedEscalationValue = false
pipelinesControllerDeployment = "tekton-pipelines-controller"

// knative/pkg StatefulSet ordinal leader-election env vars
statefulControllerOrdinalEnv = "STATEFUL_CONTROLLER_ORDINAL"
statefulReplicaCountEnv = "STATEFUL_REPLICA_COUNT"
)

// transformers that are common to all components.
Expand Down Expand Up @@ -1178,6 +1182,10 @@ func AddStatefulEnvVars(controllerName, serviceName, statefulServiceEnvVar, cont
},
},
},
{
Name: statefulReplicaCountEnv,
Value: statefulReplicaCountValue(ss),
},
}

if len(ss.Spec.Template.Spec.Containers) > 0 {
Expand All @@ -1195,6 +1203,75 @@ func AddStatefulEnvVars(controllerName, serviceName, statefulServiceEnvVar, cont
}
}

// statefulReplicaCountValue returns the effective replica count for a StatefulSet.
// Kubernetes treats a nil spec.replicas as 1.
func statefulReplicaCountValue(ss *appsv1.StatefulSet) string {
replicas := int32(1)
if ss.Spec.Replicas != nil {
replicas = *ss.Spec.Replicas
}
return strconv.Itoa(int(replicas))
}

func containerHasEnv(container corev1.Container, name string) bool {
for _, env := range container.Env {
if env.Name == name {
return true
}
}
return false
}

func upsertEnvVar(envs []corev1.EnvVar, env corev1.EnvVar) []corev1.EnvVar {
for i := range envs {
if envs[i].Name == env.Name {
envs[i] = env
return envs
}
}
return append(envs, env)
}

// SyncStatefulReplicaCountEnv updates STATEFUL_REPLICA_COUNT on StatefulSets that
// already have STATEFUL_CONTROLLER_ORDINAL, using the rendered spec.replicas.
// It does not enable ordinal mode; it only keeps the replica env in sync after
// later transformers (such as additional options) may have changed replicas.
func SyncStatefulReplicaCountEnv(manifest *mf.Manifest) error {
updated, err := manifest.Transform(func(u *unstructured.Unstructured) error {
if u.GetKind() != KindStatefulSet {
return nil
}

ss := &appsv1.StatefulSet{}
if err := runtime.DefaultUnstructuredConverter.FromUnstructured(u.Object, ss); err != nil {
return err
}
if len(ss.Spec.Template.Spec.Containers) == 0 {
return nil
}
if !containerHasEnv(ss.Spec.Template.Spec.Containers[0], statefulControllerOrdinalEnv) {
return nil
}

ss.Spec.Template.Spec.Containers[0].Env = upsertEnvVar(ss.Spec.Template.Spec.Containers[0].Env, corev1.EnvVar{
Name: statefulReplicaCountEnv,
Value: statefulReplicaCountValue(ss),
})

unstrObj, err := runtime.DefaultUnstructuredConverter.ToUnstructured(ss)
if err != nil {
return err
}
u.SetUnstructuredContent(unstrObj)
return nil
})
if err != nil {
return err
}
*manifest = updated
return nil
}

// updates performance flags/args into deployment and container given as args
// and leader election config as pod labels into a Deployment, ensuring that any changes trigger a rollout.
// It also updates the replica count if specified in the performanceSpec.
Expand Down
78 changes: 78 additions & 0 deletions pkg/reconciler/common/transformers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1137,6 +1137,84 @@ func TestAddStatefulSetPSA(t *testing.T) {
}
}

func TestAddStatefulEnvVars(t *testing.T) {
const (
controllerName = "tekton-pipelines-controller"
serviceName = "tekton-pipelines-controller"
serviceEnv = "STATEFUL_SERVICE_NAME"
ordinalEnv = "STATEFUL_CONTROLLER_ORDINAL"
)

tests := []struct {
name string
replicas *int32
wantReplicaCount string
}{
{
name: "explicit replicas",
replicas: ptr.Int32(5),
wantReplicaCount: "5",
},
{
name: "nil replicas",
replicas: nil,
wantReplicaCount: "1",
},
}

for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
ss := &appsv1.StatefulSet{
TypeMeta: metav1.TypeMeta{
Kind: KindStatefulSet,
APIVersion: appsv1.SchemeGroupVersion.String(),
},
ObjectMeta: metav1.ObjectMeta{Name: controllerName},
Spec: appsv1.StatefulSetSpec{
Replicas: test.replicas,
ServiceName: serviceName,
Template: corev1.PodTemplateSpec{
Spec: corev1.PodSpec{
Containers: []corev1.Container{{Name: "controller"}},
},
},
},
}
content, err := runtime.DefaultUnstructuredConverter.ToUnstructured(ss)
assert.NilError(t, err)
manifest, err := mf.ManifestFrom(mf.Slice([]unstructured.Unstructured{{Object: content}}))
assert.NilError(t, err)

transformed, err := manifest.Transform(AddStatefulEnvVars(controllerName, serviceName, serviceEnv, ordinalEnv))
assert.NilError(t, err)

got := statefulSetFor(t, transformed.Resources()[0])
foundService, foundOrdinal, replicaCount := false, false, ""
for _, env := range got.Spec.Template.Spec.Containers[0].Env {
switch env.Name {
case serviceEnv:
foundService = true
assert.Equal(t, env.Value, serviceName)
case ordinalEnv:
foundOrdinal = true
if env.ValueFrom == nil || env.ValueFrom.FieldRef == nil || env.ValueFrom.FieldRef.FieldPath != "metadata.name" {
t.Errorf("%s fieldPath = %v, want metadata.name", ordinalEnv, env.ValueFrom)
}
case statefulReplicaCountEnv:
replicaCount = env.Value
}
}
if !foundService {
t.Errorf("missing %s", serviceEnv)
}
if !foundOrdinal {
t.Errorf("missing %s", ordinalEnv)
}
assert.Equal(t, replicaCount, test.wantReplicaCount)
})
}
}

func TestCopyConfigMapValues(t *testing.T) {
testData := path.Join("testdata", "test-resolver-config.yaml")
manifest, err := mf.ManifestFrom(mf.Recursive(testData))
Expand Down
3 changes: 3 additions & 0 deletions pkg/reconciler/kubernetes/tektonchain/transform.go
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,9 @@ func filterAndTransform(extension common.Extension) client.FilterAndTransform {
if err := common.ExecuteAdditionalOptionsTransformer(ctx, manifest, chainCR.Spec.GetTargetNamespace(), chainCR.Spec.Options); err != nil {
return &mf.Manifest{}, err
}
if err := common.SyncStatefulReplicaCountEnv(manifest); err != nil {
return &mf.Manifest{}, err
}

return manifest, nil
}
Expand Down
11 changes: 11 additions & 0 deletions pkg/reconciler/kubernetes/tektonchain/transform_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ func TestUpdateStatefulSetOrdinalsForChains(t *testing.T) {

foundOrdinalEnv := false
foundServiceEnv := false
foundReplicaCountEnv := false
if len(sts.Spec.Template.Spec.Containers) > 0 {
for _, env := range sts.Spec.Template.Spec.Containers[0].Env {
if env.Name == tektonChainsControllerStatefulControllerOrdinal {
Expand All @@ -93,6 +94,12 @@ func TestUpdateStatefulSetOrdinalsForChains(t *testing.T) {
if env.Name == tektonChainsControllerStatefulServiceName {
foundServiceEnv = true
}
if env.Name == "STATEFUL_REPLICA_COUNT" {
foundReplicaCountEnv = true
if env.Value != "3" {
t.Errorf("Expected STATEFUL_REPLICA_COUNT to be %q, got %q", "3", env.Value)
}
}
}
}

Expand All @@ -104,6 +111,10 @@ func TestUpdateStatefulSetOrdinalsForChains(t *testing.T) {
t.Errorf("Expected to find environment variable %s", tektonChainsControllerStatefulServiceName)
}

if !foundReplicaCountEnv {
t.Errorf("Expected to find environment variable STATEFUL_REPLICA_COUNT")
}

break
}
}
Expand Down
3 changes: 3 additions & 0 deletions pkg/reconciler/kubernetes/tektonpipeline/transform.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,9 @@ func filterAndTransform(extension common.Extension) client.FilterAndTransform {
if err := common.ExecuteAdditionalOptionsTransformer(ctx, manifest, pipeline.Spec.GetTargetNamespace(), pipeline.Spec.Options); err != nil {
return &mf.Manifest{}, err
}
if err := common.SyncStatefulReplicaCountEnv(manifest); err != nil {
return &mf.Manifest{}, err
}
if pipeline.Spec.Performance.StatefulsetOrdinals != nil && *pipeline.Spec.Performance.StatefulsetOrdinals {
if err := validateStatefulSetOrdinals(manifest, leaderElectionPipelineConfig, tektonPipelinesControllerName); err != nil {
return &mf.Manifest{}, err
Expand Down
37 changes: 36 additions & 1 deletion pkg/reconciler/kubernetes/tektonpipeline/transform_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package tektonpipeline
import (
"context"
"encoding/json"
"strconv"
"testing"

"github.com/google/go-cmp/cmp"
Expand Down Expand Up @@ -344,9 +345,12 @@ func TestValidateStatefulSetOrdinalsAfterOptions(t *testing.T) {
}

manifest := ordinalManifest(t)
_, err := filterAndTransform(common.NoExtension(t.Context()))(t.Context(), &manifest, pipeline)
gotManifest, err := filterAndTransform(common.NoExtension(t.Context()))(t.Context(), &manifest, pipeline)
if test.wantError == "" {
assert.NilError(t, err)
if test.statefulSetOrdinals && test.optionReplicas != nil {
assertStatefulReplicaCountEnv(t, gotManifest, controller.statefulSet, *test.optionReplicas)
}
} else {
assert.ErrorContains(t, err, controller.configMap+test.wantError)
if test.wantStatefulSet {
Expand Down Expand Up @@ -400,6 +404,37 @@ func ordinalManifest(t *testing.T) mf.Manifest {
return manifest
}

func assertStatefulReplicaCountEnv(t *testing.T, manifest *mf.Manifest, statefulSetName string, wantReplicas int32) {
t.Helper()
found := false
for _, resource := range manifest.Resources() {
if resource.GetKind() != common.KindStatefulSet || resource.GetName() != statefulSetName {
continue
}
found = true
sts := &appsv1.StatefulSet{}
err := apimachineryRuntime.DefaultUnstructuredConverter.FromUnstructured(resource.Object, sts)
assert.NilError(t, err)
if sts.Spec.Replicas == nil || *sts.Spec.Replicas != wantReplicas {
t.Errorf("%s spec.replicas = %v, want %d", statefulSetName, sts.Spec.Replicas, wantReplicas)
}
got := ""
if len(sts.Spec.Template.Spec.Containers) > 0 {
for _, env := range sts.Spec.Template.Spec.Containers[0].Env {
if env.Name == "STATEFUL_REPLICA_COUNT" {
got = env.Value
break
}
}
}
assert.Equal(t, got, strconv.Itoa(int(wantReplicas)))
return
}
if !found {
t.Errorf("StatefulSet %s not found", statefulSetName)
}
}

// not in use, see: https://github.com/tektoncd/pipeline/pull/7789
// this field is removed from pipeline component
// keeping in types to maintain the API compatibility
Expand Down
3 changes: 3 additions & 0 deletions pkg/reconciler/kubernetes/tektonresult/transform.go
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,9 @@ func (r *Reconciler) transform(ctx context.Context, manifest *mf.Manifest, comp
if err := common.ExecuteAdditionalOptionsTransformer(ctx, manifest, instance.Spec.GetTargetNamespace(), instance.Spec.Options); err != nil {
return err
}
if err := common.SyncStatefulReplicaCountEnv(manifest); err != nil {
return err
}
return nil
}

Expand Down
11 changes: 11 additions & 0 deletions pkg/reconciler/kubernetes/tektonresult/transform_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -768,6 +768,7 @@ func TestUpdateStatefulSetOrdinalsForResults(t *testing.T) {

foundOrdinalEnv := false
foundServiceEnv := false
foundReplicaCountEnv := false
for _, container := range sts.Spec.Template.Spec.Containers {
for _, env := range container.Env {
if env.Name == tektonResultWatcherStatefulControllerOrdinal {
Expand All @@ -782,6 +783,12 @@ func TestUpdateStatefulSetOrdinalsForResults(t *testing.T) {
t.Errorf("Expected %s value to be %s, got %s", tektonResultWatcherStatefulServiceName, tektonResultWatcherServiceName, env.Value)
}
}
if env.Name == "STATEFUL_REPLICA_COUNT" {
foundReplicaCountEnv = true
if env.Value != "3" {
t.Errorf("Expected STATEFUL_REPLICA_COUNT to be %q, got %q", "3", env.Value)
}
}
}
}

Expand All @@ -793,6 +800,10 @@ func TestUpdateStatefulSetOrdinalsForResults(t *testing.T) {
t.Errorf("Expected to find environment variable %s", tektonResultWatcherStatefulServiceName)
}

if !foundReplicaCountEnv {
t.Errorf("Expected to find environment variable STATEFUL_REPLICA_COUNT")
}

break
}
}
Expand Down
Loading