Skip to content

Commit b475b1d

Browse files
committed
feat(object-controller): add standalone ClusterObjectSet manager
Add a dedicated manager with a ClusterObjectSet-only scheme and regression coverage. Reuse referenced Secret reads within each reconciliation, retry transient read failures, and show numeric revision status when no bundle version is present. Signed-off-by: Fabricio Aguiar <fabricio.aguiar@gmail.com> rh-pre-commit.version: 2.3.2 rh-pre-commit.check-secrets: ENABLED
1 parent e4d84ce commit b475b1d

7 files changed

Lines changed: 606 additions & 17 deletions

File tree

‎cmd/object-controller/main.go‎

Lines changed: 211 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,211 @@
1+
/*
2+
Copyright 2026.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package main
18+
19+
import (
20+
"crypto/tls"
21+
"flag"
22+
"fmt"
23+
"os"
24+
"time"
25+
26+
"github.com/spf13/cobra"
27+
corev1 "k8s.io/api/core/v1"
28+
"k8s.io/client-go/discovery"
29+
"k8s.io/client-go/discovery/cached/memory"
30+
_ "k8s.io/client-go/plugin/pkg/client/auth"
31+
"k8s.io/client-go/rest"
32+
"k8s.io/klog/v2"
33+
"k8s.io/utils/ptr"
34+
"pkg.package-operator.run/boxcutter/managedcache"
35+
ctrl "sigs.k8s.io/controller-runtime"
36+
"sigs.k8s.io/controller-runtime/pkg/cache"
37+
"sigs.k8s.io/controller-runtime/pkg/certwatcher"
38+
"sigs.k8s.io/controller-runtime/pkg/client"
39+
"sigs.k8s.io/controller-runtime/pkg/healthz"
40+
"sigs.k8s.io/controller-runtime/pkg/manager"
41+
"sigs.k8s.io/controller-runtime/pkg/metrics/filters"
42+
"sigs.k8s.io/controller-runtime/pkg/metrics/server"
43+
44+
ocv1 "github.com/operator-framework/operator-controller/api/v1"
45+
"github.com/operator-framework/operator-controller/internal/object-controller/controllers"
46+
"github.com/operator-framework/operator-controller/internal/object-controller/scheme"
47+
cacheutil "github.com/operator-framework/operator-controller/internal/shared/util/cache"
48+
"github.com/operator-framework/operator-controller/internal/shared/util/tlsprofiles"
49+
"github.com/operator-framework/operator-controller/internal/shared/version"
50+
)
51+
52+
type config struct {
53+
metricsAddr string
54+
pprofAddr string
55+
probeAddr string
56+
certFile string
57+
keyFile string
58+
enableLeaderElection bool
59+
}
60+
61+
func newCommand() *cobra.Command {
62+
cfg := &config{}
63+
cmd := &cobra.Command{
64+
Use: "object-controller",
65+
Short: "Manage Kubernetes objects through ClusterObjectSets",
66+
RunE: func(cmd *cobra.Command, _ []string) error {
67+
if err := cfg.validate(); err != nil {
68+
return err
69+
}
70+
restConfig, err := ctrl.GetConfig()
71+
if err != nil {
72+
return err
73+
}
74+
mgr, err := newManager(cfg, restConfig)
75+
if err != nil {
76+
return err
77+
}
78+
ctrl.Log.WithName("setup").Info("starting object-controller", "version info", version.String())
79+
return mgr.Start(cmd.Context())
80+
},
81+
}
82+
flags := cmd.Flags()
83+
flags.StringVar(&cfg.metricsAddr, "metrics-bind-address", "", "The metrics endpoint address. Disabled without tls-cert and tls-key; defaults to ':8443' when both are supplied.")
84+
flags.StringVar(&cfg.pprofAddr, "pprof-bind-address", "0", "The pprof endpoint address. An empty string or 0 disables pprof.")
85+
flags.StringVar(&cfg.probeAddr, "health-probe-bind-address", ":8081", "The health probe endpoint address.")
86+
flags.StringVar(&cfg.certFile, "tls-cert", "", "The certificate file for the metrics server. Requires tls-key.")
87+
flags.StringVar(&cfg.keyFile, "tls-key", "", "The key file for the metrics server. Requires tls-cert.")
88+
flags.BoolVar(&cfg.enableLeaderElection, "leader-elect", false, "Enable leader election for the controller manager.")
89+
logFlags := flag.NewFlagSet("logging", flag.ContinueOnError)
90+
klog.InitFlags(logFlags)
91+
flags.AddGoFlagSet(logFlags)
92+
tlsprofiles.AddFlags(flags)
93+
cmd.AddCommand(&cobra.Command{
94+
Use: "version",
95+
Short: "Print object-controller version information",
96+
Run: func(cmd *cobra.Command, _ []string) {
97+
cmd.Println(version.String())
98+
},
99+
})
100+
return cmd
101+
}
102+
103+
func (c *config) validate() error {
104+
if (c.certFile == "") != (c.keyFile == "") {
105+
return fmt.Errorf("tls-cert and tls-key flags must be used together")
106+
}
107+
if c.metricsAddr != "" && c.certFile == "" {
108+
return fmt.Errorf("metrics-bind-address requires tls-cert and tls-key flags to be set")
109+
}
110+
if c.certFile != "" && c.metricsAddr == "" {
111+
c.metricsAddr = ":8443"
112+
}
113+
return nil
114+
}
115+
116+
func newManager(cfg *config, restConfig *rest.Config) (manager.Manager, error) {
117+
metricsOptions := server.Options{BindAddress: "0"}
118+
var certWatcher *certwatcher.CertWatcher
119+
if cfg.certFile != "" {
120+
var err error
121+
certWatcher, err = certwatcher.New(cfg.certFile, cfg.keyFile)
122+
if err != nil {
123+
return nil, fmt.Errorf("initializing certificate watcher: %w", err)
124+
}
125+
tlsProfile, err := tlsprofiles.GetTLSConfigFunc()
126+
if err != nil {
127+
return nil, fmt.Errorf("getting TLS profile: %w", err)
128+
}
129+
metricsOptions = server.Options{
130+
BindAddress: cfg.metricsAddr,
131+
SecureServing: true,
132+
FilterProvider: filters.WithAuthenticationAndAuthorization,
133+
TLSOpts: []func(*tls.Config){
134+
func(c *tls.Config) {
135+
c.GetCertificate = certWatcher.GetCertificate
136+
c.NextProtos = []string{"http/1.1"}
137+
},
138+
tlsProfile,
139+
},
140+
}
141+
}
142+
mgr, err := ctrl.NewManager(restConfig, ctrl.Options{
143+
Scheme: scheme.Scheme,
144+
Metrics: metricsOptions,
145+
PprofBindAddress: cfg.pprofAddr,
146+
HealthProbeBindAddress: cfg.probeAddr,
147+
LeaderElection: cfg.enableLeaderElection,
148+
LeaderElectionID: "object-controller-lock.olm.operatorframework.io",
149+
LeaderElectionReleaseOnCancel: true,
150+
LeaseDuration: ptr.To(137 * time.Second),
151+
RenewDeadline: ptr.To(107 * time.Second),
152+
RetryPeriod: ptr.To(26 * time.Second),
153+
Cache: cache.Options{
154+
ByObject: map[client.Object]cache.ByObject{&ocv1.ClusterObjectSet{}: {}},
155+
ReaderFailOnMissingInformer: true,
156+
DefaultTransform: cacheutil.StripAnnotations(),
157+
},
158+
// References can point to immutable Secrets in any namespace. Read them
159+
// directly, without caching unrelated cluster Secrets or requiring a system namespace.
160+
Client: client.Options{Cache: &client.CacheOptions{DisableFor: []client.Object{&corev1.Secret{}}}},
161+
})
162+
if err != nil {
163+
return nil, fmt.Errorf("creating manager: %w", err)
164+
}
165+
if certWatcher != nil {
166+
if err := mgr.Add(certWatcher); err != nil {
167+
return nil, fmt.Errorf("adding certificate watcher: %w", err)
168+
}
169+
}
170+
trackingCache, err := managedcache.NewTrackingCache(
171+
ctrl.Log.WithName("trackingCache"), mgr.GetConfig(),
172+
cache.Options{Scheme: mgr.GetScheme(), Mapper: mgr.GetRESTMapper()},
173+
)
174+
if err != nil {
175+
return nil, fmt.Errorf("creating tracking cache: %w", err)
176+
}
177+
if err := mgr.Add(trackingCache); err != nil {
178+
return nil, fmt.Errorf("adding tracking cache: %w", err)
179+
}
180+
discoveryClient, err := discovery.NewDiscoveryClientForConfig(mgr.GetConfig())
181+
if err != nil {
182+
return nil, fmt.Errorf("creating discovery client: %w", err)
183+
}
184+
// Keep the field owner prefix unchanged so existing objects can be reconciled after migration.
185+
factory, err := controllers.NewDefaultRevisionEngineFactory(
186+
mgr.GetScheme(), trackingCache, memory.NewMemCacheClient(discoveryClient),
187+
mgr.GetRESTMapper(), "olm.operatorframework.io", mgr.GetConfig(),
188+
)
189+
if err != nil {
190+
return nil, fmt.Errorf("creating revision engine factory: %w", err)
191+
}
192+
if err := (&controllers.ClusterObjectSetReconciler{
193+
Client: mgr.GetClient(), RevisionEngineFactory: factory, TrackingCache: trackingCache,
194+
}).SetupWithManager(mgr); err != nil {
195+
return nil, fmt.Errorf("setting up ClusterObjectSet controller: %w", err)
196+
}
197+
if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil {
198+
return nil, err
199+
}
200+
if err := mgr.AddReadyzCheck("readyz", healthz.Ping); err != nil {
201+
return nil, err
202+
}
203+
return mgr, nil
204+
}
205+
206+
func main() {
207+
ctrl.SetLogger(klog.NewKlogr())
208+
if err := newCommand().ExecuteContext(ctrl.SetupSignalHandler()); err != nil {
209+
os.Exit(1)
210+
}
211+
}

‎cmd/object-controller/main_test.go‎

Lines changed: 164 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,164 @@
1+
package main
2+
3+
import (
4+
"bytes"
5+
"context"
6+
"encoding/json"
7+
"testing"
8+
"time"
9+
10+
"github.com/go-logr/logr/testr"
11+
"github.com/stretchr/testify/assert"
12+
"github.com/stretchr/testify/require"
13+
corev1 "k8s.io/api/core/v1"
14+
apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1"
15+
apierrors "k8s.io/apimachinery/pkg/api/errors"
16+
"k8s.io/apimachinery/pkg/api/meta"
17+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
18+
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
19+
"k8s.io/utils/ptr"
20+
ctrl "sigs.k8s.io/controller-runtime"
21+
"sigs.k8s.io/controller-runtime/pkg/client"
22+
23+
ocv1 "github.com/operator-framework/operator-controller/api/v1"
24+
"github.com/operator-framework/operator-controller/internal/object-controller/scheme"
25+
"github.com/operator-framework/operator-controller/test"
26+
)
27+
28+
func TestValidateMetricsFlags(t *testing.T) {
29+
for _, tc := range []struct {
30+
name string
31+
cfg config
32+
wantAddr string
33+
wantErr string
34+
}{
35+
{name: "disabled"},
36+
{name: "certificate without key", cfg: config{certFile: "cert"}, wantErr: "must be used together"},
37+
{name: "key without certificate", cfg: config{keyFile: "key"}, wantErr: "must be used together"},
38+
{name: "insecure metrics", cfg: config{metricsAddr: ":8443"}, wantErr: "requires tls-cert and tls-key"},
39+
{name: "default address", cfg: config{certFile: "cert", keyFile: "key"}, wantAddr: ":8443"},
40+
{name: "custom address", cfg: config{certFile: "cert", keyFile: "key", metricsAddr: ":9443"}, wantAddr: ":9443"},
41+
} {
42+
t.Run(tc.name, func(t *testing.T) {
43+
err := tc.cfg.validate()
44+
if tc.wantErr != "" {
45+
require.ErrorContains(t, err, tc.wantErr)
46+
return
47+
}
48+
require.NoError(t, err)
49+
require.Equal(t, tc.wantAddr, tc.cfg.metricsAddr)
50+
})
51+
}
52+
}
53+
54+
func TestVersionWithoutCluster(t *testing.T) {
55+
cmd := newCommand()
56+
var output bytes.Buffer
57+
cmd.SetOut(&output)
58+
cmd.SetArgs([]string{"version"})
59+
require.NoError(t, cmd.Execute())
60+
require.NotEmpty(t, output.String())
61+
}
62+
63+
// Exercise the actual manager and revision engine with only the ClusterObjectSet CRD installed.
64+
func TestStandaloneController(t *testing.T) {
65+
ctrl.SetLogger(testr.New(t))
66+
testEnv := test.NewEnv()
67+
testEnv.CRDDirectoryPaths = []string{"../../helm/olmv1/base/operator-controller/crd/experimental/olm.operatorframework.io_clusterobjectsets.yaml"}
68+
restConfig, err := testEnv.Start()
69+
require.NoError(t, err)
70+
t.Cleanup(func() { require.NoError(t, test.StopWithRetry(testEnv, time.Minute, time.Second)) })
71+
72+
cl, err := client.New(restConfig, client.Options{Scheme: scheme.Scheme})
73+
require.NoError(t, err)
74+
ctx, cancel := context.WithCancel(t.Context())
75+
defer cancel()
76+
require.True(t, apierrors.IsNotFound(cl.Get(ctx, client.ObjectKey{Name: "clusterextensions.olm.operatorframework.io"}, &apiextensionsv1.CustomResourceDefinition{})))
77+
require.True(t, apierrors.IsNotFound(cl.Get(ctx, client.ObjectKey{Name: "clustercatalogs.olm.operatorframework.io"}, &apiextensionsv1.CustomResourceDefinition{})))
78+
79+
mgr, err := newManager(&config{probeAddr: "0", pprofAddr: "0"}, restConfig)
80+
require.NoError(t, err)
81+
done := make(chan error, 1)
82+
go func() { done <- mgr.Start(ctx) }()
83+
t.Cleanup(func() {
84+
cancel()
85+
select {
86+
case err := <-done:
87+
require.NoError(t, err)
88+
case <-time.After(30 * time.Second):
89+
t.Error("manager did not stop")
90+
}
91+
})
92+
syncCtx, syncCancel := context.WithTimeout(ctx, 30*time.Second)
93+
defer syncCancel()
94+
require.True(t, mgr.GetCache().WaitForCacheSync(syncCtx), "manager cache did not synchronize")
95+
96+
for _, name := range []string{"inline", "secret-ref"} {
97+
t.Run(name, func(t *testing.T) {
98+
ns := &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{GenerateName: "standalone-"}}
99+
require.NoError(t, cl.Create(ctx, ns))
100+
manifest := unstructured.Unstructured{Object: map[string]any{
101+
"apiVersion": "v1", "kind": "ConfigMap",
102+
"metadata": map[string]any{"name": name, "namespace": ns.Name},
103+
"data": map[string]any{"hello": "world"},
104+
}}
105+
obj := ocv1.ClusterObjectSetObject{Object: manifest}
106+
if name == "secret-ref" {
107+
data, err := json.Marshal(manifest.Object)
108+
require.NoError(t, err)
109+
secret := &corev1.Secret{
110+
ObjectMeta: metav1.ObjectMeta{Name: "content", Namespace: ns.Name},
111+
Immutable: ptr.To(true), Data: map[string][]byte{"object": data},
112+
}
113+
require.NoError(t, cl.Create(ctx, secret))
114+
obj = ocv1.ClusterObjectSetObject{Ref: ocv1.ObjectSourceRef{Name: secret.Name, Namespace: secret.Namespace, Key: "object"}}
115+
}
116+
cos := &ocv1.ClusterObjectSet{
117+
ObjectMeta: metav1.ObjectMeta{Name: name},
118+
Spec: ocv1.ClusterObjectSetSpec{
119+
LifecycleState: ocv1.ClusterObjectSetLifecycleStateActive, Revision: 1,
120+
CollisionProtection: ocv1.CollisionProtectionPrevent,
121+
Phases: []ocv1.ClusterObjectSetPhase{{Name: "deploy", Objects: []ocv1.ClusterObjectSetObject{obj}}},
122+
},
123+
}
124+
require.NoError(t, cl.Create(ctx, cos))
125+
require.EventuallyWithT(t, func(collect *assert.CollectT) {
126+
if !assert.NoError(collect, cl.Get(ctx, client.ObjectKeyFromObject(cos), cos)) {
127+
return
128+
}
129+
assert.True(collect, meta.IsStatusConditionTrue(cos.Status.Conditions, ocv1.ClusterObjectSetTypeSucceeded), "%v", cos.Status.Conditions)
130+
progressing := meta.FindStatusCondition(cos.Status.Conditions, ocv1.ClusterObjectSetTypeProgressing)
131+
if assert.NotNil(collect, progressing) {
132+
assert.Equal(collect, "Revision 1 has rolled out.", progressing.Message)
133+
}
134+
}, time.Minute, 100*time.Millisecond)
135+
cm := &corev1.ConfigMap{}
136+
require.NoError(t, cl.Get(ctx, client.ObjectKey{Name: name, Namespace: ns.Name}, cm))
137+
require.Equal(t, "world", cm.Data["hello"])
138+
require.NotNil(t, metav1.GetControllerOf(cm))
139+
require.Equal(t, cos.UID, metav1.GetControllerOf(cm).UID)
140+
141+
// Observe managed-object changes without updating the ClusterObjectSet.
142+
originalUID := cm.UID
143+
require.NoError(t, cl.Delete(ctx, cm))
144+
require.EventuallyWithT(t, func(collect *assert.CollectT) {
145+
if !assert.NoError(collect, cl.Get(ctx, client.ObjectKeyFromObject(cm), cm)) {
146+
return
147+
}
148+
assert.NotEqual(collect, originalUID, cm.UID)
149+
assert.Equal(collect, "world", cm.Data["hello"])
150+
if assert.NotNil(collect, metav1.GetControllerOf(cm)) {
151+
assert.Equal(collect, cos.UID, metav1.GetControllerOf(cm).UID)
152+
}
153+
}, time.Minute, 100*time.Millisecond)
154+
155+
// The controller releases its finalizer independently of ClusterExtension.
156+
// The owner reference above lets Kubernetes garbage-collect the ConfigMap;
157+
// envtest does not run that garbage collector.
158+
require.NoError(t, cl.Delete(ctx, cos))
159+
require.Eventually(t, func() bool {
160+
return apierrors.IsNotFound(cl.Get(ctx, client.ObjectKeyFromObject(cos), cos))
161+
}, time.Minute, 100*time.Millisecond)
162+
})
163+
}
164+
}

0 commit comments

Comments
 (0)