Overview
I looked into how to write a custom controller that triggers Reconcile based on events happening outside the Kubernetes cluster, and this post writes up what I found. Once you can do this, you can build controllers that use Kubernetes custom resources to manage systems that live outside the cluster. Neat!
📘【note】 This post is day 3 of Kubernetes Advent Calendar 2021. Yesterday’s post was @makocchi’s Advanced StatefulSet を使ってみよう (Let’s try Advanced StatefulSet).
Table of Contents
- Say we want to build something like this
- controller-runtime’s Watches()
- Let’s implement it!
- The Reconcile method
- A Runnable struct that fires events periodically
- Setting up the controller
- Registering the controller with the Manager
- Summary
Say we want to build something like this
As a simple example for this post, let’s use a controller that does the following:
- Defines an object storage bucket that lives outside the Kubernetes cluster via a custom resource,
StorageBucket - Periodically health-checks the storage bucket defined by that custom resource, and records the result in the
StorageBucketresource’s Status field
To make this work, we’ll implement things so that the custom controller’s Reconcile is triggered by the event “a certain amount of time has passed.”
controller-runtime’s Watches()
controller-runtime’s builder package provides a Builder utility for controllers, which lets you build a controller that watches a given set of resources.
For example, calling builder.For() and builder.Owns() like this triggers Reconcile whenever some event happens on a ReplicaSet or Pod resource.
※ From the Controller Runtime code example
// ...(snip)...
err = builder.
ControllerManagedBy(mgr). // Create the ControllerManagedBy
For(&appsv1.ReplicaSet{}). // ReplicaSet is the Application API
Owns(&corev1.Pod{}). // ReplicaSet owns Pods created by it
Complete(&ReplicaSetReconciler{})
if err != nil {
log.Error(err, "could not create controller")
os.Exit(1)
}
// ...(snip)...
However, this only triggers Reconcile based on events happening to Kubernetes resources — it can’t trigger Reconcile based on any other kind of event, like “a certain amount of time has passed,” which is what we want for this post’s example.
For situations like that, the builder package provides the builder.Watches() method. With Watches(), you can trigger Reconcile based on an object being enqueued onto a channel. If you push an object onto that channel at regular intervals, you get exactly the behavior we’re after.
The signature of the Watches() method looks like this:
func (blder *Builder) Watches(src source.Source, eventhandler handler.EventHandler, opts ...WatchesOption) *Builder
The first and second arguments are what matter for this post, so let’s go over each of them.
- source.Source:
- The object that supplies the events that trigger Reconcile.
- source.Channel is a struct that implements this interface, wrapping a channel of
event.GenericEvent. Pushing anevent.GenericEventonto that channel is what triggers the event.
- handler.EventHandler:
- A handler that takes an event pulled off the
event.GenericEventchannel and turns it into input for Reconcile. - A GenericEvent can carry an arbitrary Kubernetes object. The EventHandler can use this to pass along a related object as input to Reconcile, or to decide whether Reconcile is even needed and skip it if not.
- A handler that takes an event pulled off the
Here’s roughly how things flow from a GenericEvent to Reconcile:
Now let’s use Watches() to implement a controller that runs Reconcile periodically.
Let’s implement it!
Here’s a walkthrough of the implementation. The code snippets quoted below are from the real thing, available in this repository — check it out alongside this post.
📘【note】 From here on, the explanation is based on the controller scaffolding generated by kubebuilder. If you haven’t used kubebuilder before, I’d recommend working through one of these tutorials first:
The Reconcile method
The Reconcile method implements the actual reconciliation logic the controller runs.
In this post’s example, we expect the StorageBucket resource’s information to arrive via the second argument, ctrl.Request.
We implement whatever check we want to run against the real storage that corresponds to the received StorageBucket resource.
func (r *StorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
log := log.FromContext(ctx)
log.Info("probe")
var storageBucket demov1.StorageBucket
if err := r.Client.Get(ctx, req.NamespacedName, &storageBucket); err != nil {
return ctrl.Result{}, client.IgnoreNotFound(err)
}
now := metav1.Now()
storageBucket.Status.LastProbeTime = &now
// Check the bucket here
storageBucket.Status.Available = true
if err := r.Client.Status().Patch(ctx, &storageBucket, client.Merge); err != nil {
return ctrl.Result{Requeue: false}, err
}
return ctrl.Result{}, nil
}
If you’ve written custom controllers before, you’ll notice this part isn’t much different from a typical one.
A Runnable struct that fires events periodically
Let’s implement a struct that fires an event once a certain amount of time has passed.
“Firing an event” here just means putting an object onto the event.GenericEvent channel.
We’ll call it ticker, and give it two methods: Start() and NeedLeaderElection().
type ticker struct {
events chan event.GenericEvent
interval time.Duration
}
// (1)
func (t *ticker) Start(ctx context.Context) error {
ticker := time.NewTicker(t.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return nil
case <-ticker.C:
t.events <- event.GenericEvent{}
}
}
}
// (2)
func (t *ticker) NeedLeaderElection() bool {
return true
}
(1)
Start()periodically pushes an event onto theevent.GenericEventchannel.- In this example we enqueue an empty
GenericEvent— all we need is a signal that time has passed, so that’s enough here.
(2)
- We’ll cover
NeedLeaderElectionbelow.
- We’ll cover
Start() and NeedLeaderElection() are the two methods defined by controller-runtime’s manager.Runnable interface (ticker is an implementation of manager.Runnable).
By implementing manager.Runnable and handing it to manager.Manager, you can let the Manager take care of safely starting and stopping the ticker, and of leader election.
Note that the Start() method’s implementation needs to follow these rules:
- It should block once called.
- It should stop blocking once the context passed as its argument is done.
NeedLeaderElection() decides whether this Runnable is subject to leader election. If you’re deploying multiple redundant Pods of your custom controller and you don’t want this Runnable running simultaneously in every Pod, have it return true.
📘【note】 @ponde_m has written a really clear post on leader election, so check that out for more detail:
Setting up the controller
The SetupWithManager method is generated by kubebuilder’s controller scaffolding, and we’ll build on that scaffolding here too.
This method creates an instance of the controller and registers it with the Manager.
We prepare the source and handler arguments that Watches() needs, and then use the builder utility — including Watches() — to create the controller instance.
We also prepare the Runnable struct that periodically fires events (more on that below).
func (c *Controller) SetupWithManager(ctx context.Context, mgr ctrl.Manager) error {
// (1)
events := make(chan event.GenericEvent)
source := source.Channel{
Source: events,
DestBufferSize: 0,
}
// (2)
err := mgr.Add(&ticker{
events: events,
interval: c.ProbeInterval,
})
if err != nil {
return err
}
// (3)
handler := handler.EnqueueRequestsFromMapFunc(func(object client.Object) []reconcile.Request {
storageBuckets := demov1.StorageBucketList{}
mgr.GetCache().List(ctx, &storageBuckets)
var requests []reconcile.Request
for _, storageBucket := range storageBuckets.Items {
requests = append(requests, reconcile.Request{
NamespacedName: types.NamespacedName{
Name: storageBucket.Name,
Namespace: storageBucket.Namespace,
},
})
}
return requests
})
// (4)
return ctrl.NewControllerManagedBy(mgr).
For(&demov1.StorageBucket{}).
Watches(&source, handler).
Complete(c)
}
- (1)
- Create the
event.GenericEventchannel.
- Create the
- (2)
- Create the Runnable struct that periodically enqueues events onto the event channel.
- The Runnable struct is registered with the Manager via
Manager.Add(), letting the Manager handle starting and stopping it (more on this below).
- (3)
- Implement a handler that turns events pulled off the event channel into Reconcile input (
reconcile.Request). - This is a somewhat blunt implementation, but it builds a
reconcile.Requestfor everyStorageBucketresource on the API server. This means Reconcile gets run against everyStorageBucketperiodically.
- Implement a handler that turns events pulled off the event channel into Reconcile input (
- (4)
- Use the builder utility to build the actual controller. The
Watches()method is givensourceandhandleras arguments.
- Use the builder utility to build the actual controller. The
Registering the controller with the Manager
From here, all that’s left is calling SetupWithManager from the main component, and you’ve got a controller that runs Reconcile on a schedule.
This part is the same as any regular custom controller, so I’ll skip the details.
That wraps up the implementation walkthrough. Nice work getting through it!
Summary
In this post, using as an example a custom controller that triggers Reconcile based on “a certain amount of time has passed,” I introduced how to build a custom controller that ties into events happening outside the Kubernetes cluster.
Let’s keep building custom controllers!