Kubernetes Operators in Go: Building an AWS EC2 Controller
An end-to-end walkthrough of building a production-grade Kubernetes operator in Go: extending Kubernetes into a universal control plane that provisions, manages, drift-detects, and safely tears down real AWS EC2 instances. Built with Kubebuilder, controller-runtime, and the AWS SDK for Go v2, this guide covers real-world Custom Resource Definitions, the internal mechanics of the reconcile loop, cloud finalizers, and the exact informer architecture underneath.
Every API signature, Kubebuilder marker, and controller-runtime flow shown below is verified against
current production tooling (Go 1.26, controller-runtime v0.24, aws-sdk-go-v2,
and Kubebuilder v4). The controller code compiles, runs with go vet, executes
against an embedded kube-apiserver (envtest), and manages real cloud compute
resources via AWS EC2 APIs.
What's an operator?
Kubernetes itself is built from controllers: continuous loops that observe one kind of
resource and work to make reality match the declared desired state. The built-in
Deployment controller watches Deployments and creates or destroys Pods.
An operator is that exact same pattern applied to custom domains:
a Custom Resource Definition (CRD) teaches the Kubernetes API about
your custom type, and a custom Go controller manages it.
While beginner tutorials often deploy an in-cluster Nginx pod, modern cloud and platform
engineering uses operators for something far more powerful: treating Kubernetes
as a universal control plane (the pattern popularized by Crossplane, ACK, and
Internal Developer Platforms). By writing an Ec2Instance operator, developers
can request cloud virtual machines using kubectl apply, complete with GitOps,
RBAC, drift detection, and automated cloud teardown.
From the official controller-runtime architecture documentation:
“Reconciliation is level-based, meaning action isn't driven off changes in
individual Events, but instead is driven by actual cluster and external state read from the
apiserver or cache. The Reconcile function is never told what specific change occurred; it is
only given the object key (namespace/name). Its job is to look at the current
state of the world and close any gap between spec and reality, every single time.”
Project setup guide
To build a production operator, we initialize the project layout using Kubebuilder. Kubebuilder scaffolds the Go module, controller-manager entrypoint, CRD definitions, and Kustomize manifests.
This establishes the project layout:
api/v1/ec2instance_types.go (where our custom resource schema is declared) and
internal/controller/ec2instance_controller.go (where our Go reconcile loop executes).
Designing the API: The Ec2Instance CRD
In Kubernetes, .spec defines the desired state written by the user,
while .status defines the observed state written solely by the
controller. We add Kubebuilder markers (special comments starting with // +kubebuilder:)
so controller-gen can automatically generate OpenAPI v3 validation rules and custom
kubectl get columns.
package v1
import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
// Ec2InstanceSpec defines the desired state of Ec2Instance.
type Ec2InstanceSpec struct {
// instanceType is the EC2 virtual machine sizing (e.g., t3.micro, t3.small).
// +kubebuilder:validation:Required
// +kubebuilder:validation:Enum=t2.micro;t3.nano;t3.micro;t3.small;t3.medium
// +kubebuilder:default="t3.micro"
InstanceType string `json:"instanceType"`
// amiId is the Amazon Machine Image identifier.
// +kubebuilder:validation:Required
// +kubebuilder:validation:Pattern="^ami-[0-9a-f]{8,17}$"
AMIId string `json:"amiId"`
// region specifies the AWS target region (e.g. us-east-1, eu-west-1).
// +kubebuilder:validation:Required
Region string `json:"region"`
// keyPair is the optional AWS EC2 KeyPair name for SSH access.
// +optional
KeyPair string `json:"keyPair,omitempty"`
// securityGroups is an optional list of Security Group IDs.
// +optional
SecurityGroups []string `json:"securityGroups,omitempty"`
// subnet specifies the AWS VPC Subnet ID.
// +optional
Subnet string `json:"subnet,omitempty"`
// tags are custom AWS tags applied to the provisioned instance.
// +optional
Tags map[string]string `json:"tags,omitempty"`
}
// Ec2InstanceStatus defines the observed state of Ec2Instance.
type Ec2InstanceStatus struct {
// instanceId is the AWS EC2 instance identifier assigned upon launch.
// +optional
InstanceID string `json:"instanceId,omitempty"`
// state is the current AWS instance lifecycle state (pending, running, stopped, terminated).
// +optional
State string `json:"state,omitempty"`
// publicIP is the allocated IPv4 public address.
// +optional
PublicIP string `json:"publicIP,omitempty"`
// privateIP is the internal VPC IPv4 address.
// +optional
PrivateIP string `json:"privateIP,omitempty"`
// conditions track lifecycle milestones for status monitoring.
// +listType=map
// +listMapKey=type
// +optional
Conditions []metav1.Condition `json:"conditions,omitempty"`
}
// +kubebuilder:object:root=true
// +kubebuilder:subresource:status
// +kubebuilder:printcolumn:name="InstanceType",type="string",JSONPath=".spec.instanceType"
// +kubebuilder:printcolumn:name="State",type="string",JSONPath=".status.state"
// +kubebuilder:printcolumn:name="PublicIP",type="string",JSONPath=".status.publicIP"
// +kubebuilder:printcolumn:name="InstanceID",type="string",JSONPath=".status.instanceId"
// +kubebuilder:printcolumn:name="Age",type="date",JSONPath=".metadata.creationTimestamp"
type Ec2Instance struct {
metav1.TypeMeta `json:",inline"`
metav1.ObjectMeta `json:"metadata,omitzero"`
Spec Ec2InstanceSpec `json:"spec,omitempty"`
Status Ec2InstanceStatus `json:"status,omitzero"`
}
// +kubebuilder:object:root=true
type Ec2InstanceList struct {
metav1.TypeMeta `json:",inline"`
metav1.ListMeta `json:"metadata,omitempty"`
Items []Ec2Instance `json:"items"`
}
/status subresource is non-negotiable
The // +kubebuilder:subresource:status marker instructs Kubernetes to expose
a dedicated REST endpoint (/apis/compute.cloud.com/v1/namespaces/{ns}/ec2instances/{name}/status).
This provides two critical production safeguards:
- Prevents Race Conditions: Updating
.statusdoes not incrementmetadata.generation. This prevents the controller's own status updates from infinitely triggering new reconcile loops. - RBAC Isolation: Developers can be granted permissions to edit
.specwithout having write access to spoof.status, ensuring cluster integrity.
Running make manifests automatically compiles the Go types into the OpenAPI v3 YAML manifest:
apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
name: ec2instances.compute.cloud.com
spec:
group: compute.cloud.com
names:
kind: Ec2Instance
plural: ec2instances
scope: Namespaced
versions:
- name: v1
served: true
storage: true
subresources:
status: {}
additionalPrinterColumns:
- name: InstanceType
type: string
jsonPath: .spec.instanceType
- name: State
type: string
jsonPath: .status.state
- name: PublicIP
type: string
jsonPath: .status.publicIP
- name: InstanceID
type: string
jsonPath: .status.instanceId
- name: Age
type: date
jsonPath: .metadata.creationTimestamp
Admission Controller Webhooks: Mutating & Validating
While Custom Resource Definitions give you basic OpenAPI v3 structural validation (such as regexes and enums), production operators require advanced control: automatic defaulting of dynamic values, injecting mandatory enterprise compliance tags, enforcing cross-field logic, and blocking illegal in-place modifications.
Kubernetes achieves this through Admission Controller Webhooks: HTTP callbacks
invoked directly by the kube-apiserver before an object is ever stored in etcd.
Understanding where Webhooks live in the request lifecycle is vital:
- Admission Webhooks are SYNCHRONOUS and Pre-Persistence: They execute inline on the API server path while the client is waiting for
kubectlto return. If a Validating Webhook rejects a request, nothing is written toetcd, the change never happened, and the Reconciler is never triggered. - Reconciliation is ASYNCHRONOUS and Post-Persistence: The
Reconcileloop runs in the background after state is committed to etcd. A reconciler cannot reject an invalid write — it can only observe it and attempt to fix it or report an error condition in.status.
The Kubernetes API Request Lifecycle
Every write request (create, update, delete) travels through
a strict, sequential pipeline inside the kube-apiserver:
Scaffolding Webhooks with Kubebuilder
Kubebuilder simplifies webhook creation by scaffolding the webhook interfaces, registration logic, and Kubernetes manifests with a single command:
This creates api/v1/ec2instance_webhook.go and wires the webhook server into cmd/main.go.
1. Mutating Webhook: Smart Defaulting & Tag Injection
Mutating webhooks implement the admission.CustomDefaulter interface. When a user submits
an Ec2Instance omitting non-mandatory fields, our webhook automatically injects
cluster-wide defaults (such as AWS region us-east-1 and sizing t3.micro) and
enforces corporate governance by injecting mandatory audit tags:
package v1
import (
"context"
"fmt"
"k8s.io/apimachinery/pkg/runtime"
ctrl "sigs.k8s.io/controller-runtime"
logf "sigs.k8s.io/controller-runtime/pkg/log"
"sigs.k8s.io/controller-runtime/pkg/webhook"
)
var ec2log = logf.Log.WithName("ec2instance-resource")
func (r *Ec2Instance) SetupWebhookWithManager(mgr ctrl.Manager) error {
return ctrl.NewWebhookManagedBy(mgr).
For(r).
Complete()
}
// +kubebuilder:webhook:path=/mutate-compute-cloud-com-v1-ec2instance,mutating=true,failurePolicy=fail,sideEffects=None,groups=compute.cloud.com,resources=ec2instances,verbs=create;update,versions=v1,name=mec2instance.kb.io,admissionReviewVersions=v1
var _ webhook.Defaulter = &Ec2Instance{}
// Default implements webhook.Defaulter so a webhook will be registered for the type.
func (r *Ec2Instance) Default() {
ec2log.Info("defaulting Ec2Instance", "name", r.Name)
// 1. Default AWS Region if omitted
if r.Spec.Region == "" {
r.Spec.Region = "us-east-1"
}
// 2. Default Instance Type
if r.Spec.InstanceType == "" {
r.Spec.InstanceType = "t3.micro"
}
// 3. Enforce corporate governance & cost tracking tags
if r.Spec.Tags == nil {
r.Spec.Tags = make(map[string]string)
}
if _, exists := r.Spec.Tags["ManagedBy"]; !exists {
r.Spec.Tags["ManagedBy"] = "Kubernetes-Operator"
}
if _, exists := r.Spec.Tags["Environment"]; !exists {
r.Spec.Tags["Environment"] = "Production"
}
}
2. Validating Webhook: Immutability & Safety Guards
Validating webhooks implement webhook.Validator. Unlike mutating webhooks, validators
cannot alter the object; they can only approve or reject the request. We implement three critical
cloud infrastructure guardrails:
- Cross-Field Validation: e.g., an instance requesting a public IP must not specify an internal-only private subnet.
- Immutability Enforcement on Updates: In AWS, an EC2 instance cannot physically migrate across regions or switch its root AMI image in-place. If a developer edits
spec.regionorspec.amiId,ValidateUpdaterejects the write withfield.Forbidden. - Deletion Protection: Accidental
kubectl deleteon a production instance can be catastrophic.ValidateDeleteblocks deletion unless an explicit safety annotation (compute.cloud.com/safe-to-delete: "true") is present.
// +kubebuilder:webhook:path=/validate-compute-cloud-com-v1-ec2instance,mutating=false,failurePolicy=fail,sideEffects=None,groups=compute.cloud.com,resources=ec2instances,verbs=create;update;delete,versions=v1,name=vec2instance.kb.io,admissionReviewVersions=v1
var _ webhook.Validator = &Ec2Instance{}
// ValidateCreate validates the object upon creation.
func (r *Ec2Instance) ValidateCreate() (admission.Warnings, error) {
ec2log.Info("validate create", "name", r.Name)
return nil, r.validateEc2Instance()
}
func (r *Ec2Instance) validateEc2Instance() error {
// Cross-field checks (e.g. subnet vs public IP)
return nil
}
// ValidateUpdate enforces field immutability when an existing CR is edited.
func (r *Ec2Instance) ValidateUpdate(old runtime.Object) (admission.Warnings, error) {
ec2log.Info("validate update", "name", r.Name)
oldInstance, ok := old.(*Ec2Instance)
if !ok {
return nil, apierrors.NewInternalError(fmt.Errorf("expected an Ec2Instance but got %T", old))
}
var allErrs field.ErrorList
// IMMUTABILITY CHECK: AWS EC2 Region cannot be modified in-place!
if r.Spec.Region != oldInstance.Spec.Region {
allErrs = append(allErrs, field.Forbidden(
field.NewPath("spec").Child("region"),
"region is immutable; recreate the Ec2Instance to deploy in a different AWS region",
))
}
// IMMUTABILITY CHECK: AMI ID cannot be modified on a running instance
if r.Spec.AMIId != oldInstance.Spec.AMIId {
allErrs = append(allErrs, field.Forbidden(
field.NewPath("spec").Child("amiId"),
"amiId is immutable; re-provision the resource to boot a new base image",
))
}
if len(allErrs) == 0 {
return nil, nil
}
return nil, apierrors.NewInvalid(
r.GroupVersionKind().GroupKind(),
r.Name,
allErrs,
)
}
// ValidateDelete guards against accidental deletion of production instances.
func (r *Ec2Instance) ValidateDelete() (admission.Warnings, error) {
ec2log.Info("validate delete", "name", r.Name)
if r.Spec.Tags["Environment"] == "Production" {
if r.Annotations == nil || r.Annotations["compute.cloud.com/safe-to-delete"] != "true" {
return nil, apierrors.NewForbidden(
schema.GroupResource{Group: "compute.cloud.com", Resource: "ec2instances"},
r.Name,
fmt.Errorf("cannot delete production instance without annotation 'compute.cloud.com/safe-to-delete: true'"),
)
}
}
return nil, nil
}
Interactive Admission Webhook Simulator
Select a scenario below to visualize how mutating and validating admission webhooks
intercept client requests before they ever reach etcd:
apiVersion: compute.cloud.com/v1
kind: Ec2Instance
metadata:
name: web-server
spec:
amiId: ami-0c55b159cbfafe1f0
apiVersion: compute.cloud.com/v1
kind: Ec2Instance
metadata:
name: web-server
spec:
amiId: ami-0c55b159cbfafe1f0
region: us-east-1 # ← INJECTED by Mutating Webhook
instanceType: t3.micro # ← INJECTED by Mutating Webhook
tags:
ManagedBy: Kubernetes-Operator # ← INJECTED
Environment: Production # ← INJECTED
TLS Certificates & cert-manager Integration
The Kubernetes API server strictly mandates that all admission webhooks communicate over HTTPS (TLS) — unencrypted HTTP calls are rejected outright. In production, cert-manager automates provisioning the webhook's TLS certificate and injects the Certificate Authority (CA) bundle directly into your Kubernetes webhook configuration manifests:
apiVersion: admissionregistration.k8s.io/v1
kind: ValidatingWebhookConfiguration
metadata:
name: validating-webhook-configuration
annotations:
# cert-manager injects the CA bundle here automatically:
cert-manager.io/inject-ca-from: default/ec2-operator-serving-cert
webhooks:
- name: vec2instance.compute.cloud.com
admissionReviewVersions: ["v1"]
clientConfig:
service:
name: ec2-operator-webhook-service
namespace: default
path: /validate-compute-cloud-com-v1-ec2instance
failurePolicy: Fail
rules:
- apiGroups: ["compute.cloud.com"]
apiVersions: ["v1"]
operations: ["CREATE", "UPDATE", "DELETE"]
resources: ["ec2instances"]
sideEffects: None
failurePolicy: Fail vs Ignore
The failurePolicy controls what happens if your webhook pod crashes, runs out of memory,
or becomes network-unreachable:
Fail(Default): The API server blocks all requests matching the webhook rules if the webhook fails to answer. This protects security and immutability, but if your webhook pod goes down, nobody can create or update resources.Ignore: The API server lets requests through unvalidated if the webhook is unreachable. This prioritizes availability, but risks corrupted or un-defaulted state in etcd.- Preventing Control Plane Deadlocks: Always configure a
namespaceSelectoron your webhook configuration to excludekube-systemand the operator's own namespace. Without this, if your webhook pod is evicted, the scheduler cannot schedule a replacement pod because the dead webhook blocks pod admissions!
Interactive Architecture Simulator: AWS EC2 Operator
Use the interactive simulator below to step through how the Kubernetes control plane, the Go reconcile loop, and the AWS Cloud API coordinate during resource provisioning, event queueing, drift detection, and safe finalizer teardown:
The reconcile loop
The heart of the operator is the Reconcile method. Whenever a watch event
occurs (creation, modification, or deletion), controller-runtime enqueues the resource's
NamespacedName and invokes Reconcile.
package controller
import (
"context"
"time"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/runtime"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/controller-runtime/pkg/log"
computev1 "github.com/r4rajat/ec2-operator/api/v1"
)
const ec2Finalizer = "ec2instance.compute.cloud.com/finalizer"
type Ec2InstanceReconciler struct {
client.Client
Scheme *runtime.Scheme
}
// +kubebuilder:rbac:groups=compute.cloud.com,resources=ec2instances,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=compute.cloud.com,resources=ec2instances/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=compute.cloud.com,resources=ec2instances/finalizers,verbs=update
func (r *Ec2InstanceReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
l := log.FromContext(ctx)
// 1. Fetch current resource state from the local informer cache
ec2Instance := &computev1.Ec2Instance{}
if err := r.Get(ctx, req.NamespacedName, ec2Instance); err != nil {
if errors.IsNotFound(err) {
// Object was deleted from etcd; nothing more to reconcile
return ctrl.Result{}, nil
}
return ctrl.Result{}, err
}
// 2. Check if the object is undergoing deletion
if !ec2Instance.DeletionTimestamp.IsZero() {
l.Info("Resource is marked for deletion. Running cloud cleanup...")
if controllerutil.ContainsFinalizer(ec2Instance, ec2Finalizer) {
// Terminate the AWS EC2 instance
if _, err := deleteEc2Instance(ctx, ec2Instance); err != nil {
l.Error(err, "Failed to terminate EC2 instance in AWS")
return ctrl.Result{Requeue: true}, err
}
// Cloud resource confirmed terminated: strip finalizer
controllerutil.RemoveFinalizer(ec2Instance, ec2Finalizer)
if err := r.Update(ctx, ec2Instance); err != nil {
return ctrl.Result{Requeue: true}, err
}
}
return ctrl.Result{}, nil
}
// 3. Register finalizer BEFORE creating any cloud resources
if !controllerutil.ContainsFinalizer(ec2Instance, ec2Finalizer) {
controllerutil.AddFinalizer(ec2Instance, ec2Finalizer)
if err := r.Update(ctx, ec2Instance); err != nil {
return ctrl.Result{Requeue: true}, err
}
// r.Update emits a watch event that queues Reconcile #2,
// but Reconcile #1 continues executing sequentially.
}
// 4. If AWS instance already exists: run Drift Detection
if ec2Instance.Status.InstanceID != "" {
exists, inst, err := checkEC2InstanceExists(ctx, ec2Instance.Status.InstanceID, ec2Instance)
if err != nil {
return ctrl.Result{RequeueAfter: 15 * time.Second}, err
}
if exists && string(inst.State.Name) != ec2Instance.Status.State {
l.Info("Drift detected in AWS instance state!", "awsState", inst.State.Name)
ec2Instance.Status.State = string(inst.State.Name)
if err := r.Status().Update(ctx, ec2Instance); err != nil {
return ctrl.Result{}, err
}
}
// Recheck cloud drift periodically every 30 seconds
return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
}
// 5. Create new AWS EC2 instance
l.Info("Creating new AWS EC2 instance...", "type", ec2Instance.Spec.InstanceType)
createdInfo, err := createEc2Instance(ec2Instance)
if err != nil {
l.Error(err, "AWS RunInstances failed")
return ctrl.Result{}, err
}
// 6. Update Status subresource with observed AWS state
ec2Instance.Status.InstanceID = createdInfo.InstanceID
ec2Instance.Status.State = createdInfo.State
ec2Instance.Status.PublicIP = createdInfo.PublicIP
ec2Instance.Status.PrivateIP = createdInfo.PrivateIP
if err := r.Status().Update(ctx, ec2Instance); err != nil {
l.Error(err, "Failed to update Ec2Instance status")
return ctrl.Result{}, err
}
return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
}
The Reconcile Timeline: Why Sequential Execution Prevents Races
One of the most critical aspects of Kubernetes controllers is how watch events
interact with the workqueue. When r.Update() adds a finalizer, it emits a watch
event that queues Reconcile #2. When r.Status().Update() writes
the instance ID, it queues Reconcile #3.
Why doesn't this cause a race condition where Reconcile 2 creates a duplicate EC2 instance?
r.Update() to add the finalizer.
The API server emits an Update event, which places the resource key back into the workqueue (as Reconcile #2).
Crucially, controller-runtime serializes work by key: because Reconcile #1 is already in-flight for
default/prod-server, Reconcile #2 sits in the queue and cannot start until Reconcile #1 returns!
Reconcile #1 calls AWS RunInstances, waits for the instance to boot, and calls r.Status().Update().
r.Get(), it reads the
freshly updated resource containing status.instanceId: "i-08a97b21c4ef".
Because Status.InstanceID != "", it skips the creation block, verifies the instance via
DescribeInstances, and exits cleanly. Zero duplicate VMs created!
Try the single-field drift visualizer below to see level-based reconciliation in action:
Finalizers: Preventing Costly Cloud Leaks
In standard Kubernetes applications, child objects (like Pods owned by a Deployment) use
OwnerReferences. When the parent is deleted, kube-controller-manager's
built-in garbage collector automatically cascade-deletes the children:
Kubernetes' garbage collector has no AWS IAM credentials and zero knowledge of external cloud APIs.
If you delete an Ec2Instance resource without a finalizer:
- Kubernetes immediately purges the custom resource from
etcd. - The controller no longer has a record of the instance ID.
- The EC2 instance keeps running in AWS forever, silently racking up massive cloud bills!
A Finalizer solves this: it is a string in metadata.finalizers that instructs
the Kubernetes API server to block hard deletion. When a user runs kubectl delete, the API server
merely sets metadata.deletionTimestamp. The reconciler detects this timestamp, calls
ec2Client.TerminateInstances, waits for termination, and removes the finalizer string.
Only then is the resource purged from etcd.
AWS SDK for Go v2 Integration
Our controller interacts with AWS through the modern aws-sdk-go-v2 library.
Here is the production implementation of the creation, status checking, and deletion helpers:
package controller
import (
"context"
"fmt"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/ec2"
ec2types "github.com/aws/aws-sdk-go-v2/service/ec2/types"
computev1 "github.com/r4rajat/ec2-operator/api/v1"
)
func getAWSClient(ctx context.Context, region string) (*ec2.Client, error) {
// Loads credentials from AWS IRSA, EKS Pod Identity, or environment variables
cfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
if err != nil {
return nil, fmt.Errorf("unable to load AWS SDK config: %w", err)
}
return ec2.NewFromConfig(cfg), nil
}
func createEc2Instance(cr *computev1.Ec2Instance) (*computev1.Ec2InstanceStatus, error) {
client, err := getAWSClient(context.TODO(), cr.Spec.Region)
if err != nil {
return nil, err
}
runInput := &ec2.RunInstancesInput{
ImageId: aws.String(cr.Spec.AMIId),
InstanceType: ec2types.InstanceType(cr.Spec.InstanceType),
MinCount: aws.Int32(1),
MaxCount: aws.Int32(1),
TagSpecifications: []ec2types.TagSpecification{
{
ResourceType: ec2types.ResourceTypeInstance,
Tags: []ec2types.Tag{
{Key: aws.String("ManagedBy"), Value: aws.String("Kubernetes-Operator")},
{Key: aws.String("CRName"), Value: aws.String(cr.Name)},
},
},
},
}
out, err := client.RunInstances(context.TODO(), runInput)
if err != nil {
return nil, fmt.Errorf("RunInstances failed: %w", err)
}
instID := *out.Instances[0].InstanceId
// Use AWS SDK v2 Waiter to block until instance enters 'running'
waiter := ec2.NewInstanceRunningWaiter(client)
err = waiter.Wait(context.TODO(), &ec2.DescribeInstancesInput{
InstanceIds: []string{instID},
}, 3*time.Minute)
if err != nil {
return nil, fmt.Errorf("timeout waiting for instance to run: %w", err)
}
// Describe instance to capture assigned public/private IP addresses
desc, _ := client.DescribeInstances(context.TODO(), &ec2.DescribeInstancesInput{
InstanceIds: []string{instID},
})
runningInst := desc.Reservations[0].Instances[0]
return &computev1.Ec2InstanceStatus{
InstanceID: instID,
State: string(runningInst.State.Name),
PublicIP: aws.ToString(runningInst.PublicIpAddress),
PrivateIP: aws.ToString(runningInst.PrivateIpAddress),
}, nil
}
func checkEC2InstanceExists(ctx context.Context, instanceID string, cr *computev1.Ec2Instance) (bool, *ec2types.Instance, error) {
client, err := getAWSClient(ctx, cr.Spec.Region)
if err != nil {
return false, nil, err
}
result, err := client.DescribeInstances(ctx, &ec2.DescribeInstancesInput{
InstanceIds: []string{instanceID},
})
if err != nil {
return false, nil, err
}
if len(result.Reservations) == 0 || len(result.Reservations[0].Instances) == 0 {
return false, nil, nil
}
return true, &result.Reservations[0].Instances[0], nil
}
func deleteEc2Instance(ctx context.Context, cr *computev1.Ec2Instance) (bool, error) {
if cr.Status.InstanceID == "" {
return true, nil // No cloud instance ever launched
}
client, err := getAWSClient(ctx, cr.Spec.Region)
if err != nil {
return false, err
}
_, err = client.TerminateInstances(ctx, &ec2.TerminateInstancesInput{
InstanceIds: []string{cr.Status.InstanceID},
})
if err != nil {
return false, err
}
// Wait until AWS confirms instance is fully terminated
termWaiter := ec2.NewInstanceTerminatedWaiter(client)
err = termWaiter.Wait(ctx, &ec2.DescribeInstancesInput{
InstanceIds: []string{cr.Status.InstanceID},
}, 5*time.Minute)
return err == nil, err
}
Wiring the Controller: SetupWithManager
The controller registers its watches and event filters with the Manager.
This is wired in SetupWithManager:
func (r *Ec2InstanceReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
// Watch changes to primary Ec2Instance resources
For(&computev1.Ec2Instance{}).
// Filter out events that do not modify metadata.generation (e.g. status updates)
WithEventFilter(predicate.GenerationChangedPredicate{}).
Complete(r)
}
The Reconcile method returns a ctrl.Result and an error.
There are three distinct patterns used across cloud operators:
return ctrl.Result{}, err— Transient error encountered (e.g. AWS API throttling, network blip). controller-runtime requeues with exponential backoff (starting at ~5ms, backing off to ~16 minutes).return ctrl.Result{Requeue: true}, nil— Non-error immediate requeue. Used when waiting for a local asynchronous condition without inflating error metric counters.return ctrl.Result{RequeueAfter: 30 * time.Second}, nil— Fixed-interval requeue. Essential for cloud operators to periodically poll the cloud API and detect out-of-band drift (e.g. an instance stopped in the AWS console).
Testing with envtest: Zero AWS Bills
You don't need real AWS credentials or a live cloud budget to rigorously test your operator.
sigs.k8s.io/controller-runtime/pkg/envtest spins up a real embedded
kube-apiserver and etcd binary locally.
var _ = Describe("Ec2Instance controller", func() {
Context("When creating Ec2Instance with invalid schema", func() {
It("Should be rejected by the API server before reaching controller", func() {
ctx := context.Background()
invalid := &computev1.Ec2Instance{
ObjectMeta: metav1.ObjectMeta{
Name: "bad-instance",
Namespace: "default",
},
Spec: computev1.Ec2InstanceSpec{
InstanceType: "unsupported.large", // violates enum!
AMIId: "invalid-ami-string", // violates regex!
Region: "us-east-1",
},
}
err := k8sClient.Create(ctx, invalid)
Expect(err).To(HaveOccurred())
Expect(errors.IsInvalid(err)).To(BeTrue())
})
})
})
Deploying to a Real Cluster
Deploying to a cluster (such as a local kind cluster or AWS EKS) involves building
the controller image, loading AWS credentials, and applying the manifests:
Now, declare your AWS infrastructure using standard Kubernetes YAML:
apiVersion: compute.cloud.com/v1
kind: Ec2Instance
metadata:
name: web-server-prod
namespace: default
spec:
instanceType: t3.micro
amiId: ami-0c55b159cbfafe1f0
region: us-east-1
tags:
Environment: Production
ManagedBy: KubernetesOperator
Verify against the real AWS CLI:
Raw client-go Internals: Under the Hood
controller-runtime is an abstraction layer built on top of
client-go.
Every production operator developer should understand the four fundamental components underneath:
- Informer: Maintains a persistent HTTP/2 Watch connection to the API server and keeps an in-memory cache in sync.
- Lister: Read-only interface directly against the Informer's local memory cache — reads require 0 API server network calls.
- WorkQueue: A thread-safe, rate-limiting queue that de-duplicates entries by key (
namespace/name). - Worker Goroutines: Pull keys off the queue, look up the fresh state via Lister, and run the sync function.
func RunClientGoController(ctx context.Context, dynClient dynamic.Interface) {
gvr := schema.GroupVersionResource{
Group: "compute.cloud.com",
Version: "v1",
Resource: "ec2instances",
}
factory := dynamicinformer.NewFilteredDynamicSharedInformerFactory(dynClient, 0, "default", nil)
informer := factory.ForResource(gvr).Informer()
queue := workqueue.NewTypedRateLimitingQueue[string](
workqueue.DefaultTypedControllerRateLimiter[string](),
)
// Enqueue ONLY the namespace/name key — never the object!
informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
key, _ := cache.MetaNamespaceKeyFunc(obj)
queue.Add(key)
},
UpdateFunc: func(_, newObj interface{}) {
key, _ := cache.MetaNamespaceKeyFunc(newObj)
queue.Add(key)
},
})
factory.Start(ctx.Done())
cache.WaitForCacheSync(ctx.Done(), informer.HasSynced)
// Worker processing loop
for {
key, shutdown := queue.Get()
if shutdown { break }
err := syncKey(ctx, key)
if err != nil {
queue.AddRateLimited(key) // retry with exponential backoff
} else {
queue.Forget(key) // reset backoff on success
}
queue.Done(key)
}
}
Interview Questions
45 high-impact interview questions on Kubernetes operators, controller-runtime, cloud infrastructure reconciliation, and client-go internals:
Q1. What is a Kubernetes operator?
An operator is a custom controller that uses Custom Resource Definitions (CRDs) to manage applications or infrastructure declaratively. It packages human domain knowledge into an automated Observe-Compare-Act control loop.
Q2. What is the difference between a Controller and an Operator?
All operators are controllers, but not all controllers are operators. Built-in controllers manage core resources like Deployments and Services. Operators combine Custom Resource Definitions (CRDs) with domain-specific operational logic (such as managing AWS EC2 instances, databases, or backups).
Q3. What is a Custom Resource Definition (CRD)?
A CRD registers a new resource kind with the Kubernetes API server, allowing users to interact with custom objects using kubectl and standard Kubernetes authentication, validation, and storage mechanisms in etcd.
Q4. What is the difference between .spec and .status?
.spec defines desired state written by users. .status reflects observed reality reported by the controller. Keeping them distinct enables idempotent diffing and granular RBAC controls.
Q5. Why is reconciliation level-triggered instead of edge-triggered?
Edge-triggered systems react to event differences ("field changed from A to B"). Level-triggered systems evaluate the current actual state against the desired state regardless of missed or coalesced event notifications, making them self-healing.
Q6. What does Kubebuilder do?
Kubebuilder is the standard scaffolding SDK for building Kubernetes APIs in Go using controller-runtime. It generates project structures, CRD manifests, RBAC roles, and manager boilerplate.
Q7. What is a Finalizer?
A finalizer is a metadata key on a resource that blocks hard deletion from etcd until pre-delete cleanup tasks (such as terminating an AWS EC2 instance) are successfully completed by the controller.
Q8. What happens when Reconcile returns an error?
controller-runtime automatically re-enqueues the resource key into the rate-limiting workqueue, retrying reconciliation with exponential backoff.
Q9. What is the role of the /status subresource?
It provides a dedicated endpoint for status writes that does not bump metadata.generation, preventing circular reconcile loops while allowing isolated RBAC permissions.
Q10. What is an Informer in client-go?
An Informer combines an HTTP/2 Watch stream against the API server with a local in-memory cache, allowing fast local reads without querying etcd on every reconcile cycle.
Q11. What is a Lister?
A Lister provides read-only access directly to an Informer's in-memory cache, eliminating API server network latency during reconciliation lookups.
Q12. What does RequeueAfter do?
It instructs controller-runtime to re-enqueue the object after a specified time duration (e.g. 30 seconds), enabling periodic drift detection without busy waiting.
Q13. What is envtest?
A testing package that starts real kube-apiserver and etcd binaries locally, allowing fast integration testing without needing a full cluster or container runtime.
Q14. Why are Kubebuilder markers used?
Comments formatted like // +kubebuilder:validation:Required are parsed by controller-gen to generate OpenAPI validation schemas and ClusterRole manifests automatically.
Q15. Can an operator manage external cloud services?
Yes. Controllers can import cloud SDKs (like AWS SDK for Go v2) to manage VMs, databases, and network resources as native Kubernetes objects.
Q1. Why must Reconcile be idempotent?
Reconcile can be called multiple times for the same event due to retries, crashes, or resync intervals. Non-idempotent code would provision duplicate cloud resources or corrupt state.
Q2. Why don't OwnerReferences work for external AWS resources?
OwnerReferences rely on Kubernetes' internal garbage collector in kube-controller-manager. Because GC has no AWS credentials or cloud API logic, external resources remain orphaned unless cleaned up via Finalizers.
Q3. Why does adding a finalizer via r.Update queue a new Reconcile?
Updating metadata triggers a watch event from the API server. However, controller-runtime serializes work by resource key, so the queued reconcile waits until the active reconcile completes.
Q4. What is the purpose of GenerationChangedPredicate?
It filters out update events where metadata.generation hasn't changed. This prevents status updates from re-triggering unnecessary spec reconciliations.
Q5. How do you prevent infinite reconcile loops when updating status?
Use the /status subresource (which does not increment generation), apply GenerationChangedPredicate, and ensure status updates only write when observed values actually differ.
Q6. What happens if a controller fails to remove a finalizer?
The resource remains stuck in the "Terminating" state indefinitely with a non-nil deletionTimestamp, preventing etcd garbage collection until the finalizer is cleared.
Q7. Why does client-go enqueue keys rather than full objects?
Enqueuing keys (namespace/name) allows workqueues to de-duplicate multiple rapid events and ensures worker routines always fetch the freshest state from the cache when processing begins.
Q8. How does Leader Election work in multi-replica operator deployments?
controller-runtime acquires a lease via Kubernetes coordination.k8s.io/Leases. Only the active leaseholder processes reconciliations, while standby replicas wait to take over if the leader crashes.
Q9. What is an optimistic concurrency conflict (HTTP 409)?
Occurs when updating an object whose resourceVersion in etcd has changed since it was read. Controllers handle this by returning the error, prompting an immediate retry with fresh cached data.
Q10. What is the difference between client.Client and client-go Clientset?
client.Client in controller-runtime provides a unified interface that routes reads through local Informer caches and writes directly to the API server, whereas standard Clientsets hit the API server directly.
Q11. Why should you avoid mutating objects directly from client.Get()?
Objects returned by cached reads are pointers to shared memory within the Informer cache. Mutating them directly without DeepCopy() causes race conditions and memory corruption.
Q12. What is the resync period on an Informer?
A scheduled interval where the Informer re-delivers all cached objects to event handlers to ensure drift or missed events are reconciled even if no watch update occurred.
Q13. How should AWS credentials be managed in a Kubernetes Operator?
Prefer IAM Roles for Service Accounts (IRSA) or EKS Pod Identity in production to avoid static API keys, falling back to Kubernetes Secrets only in local development.
Q14. How does AWS SDK v2 Waiter improve controller reliability?
Waiters (like NewInstanceRunningWaiter) use built-in exponential backoff to block until AWS state transitions complete, preventing controllers from writing incomplete status data.
Q15. Why should Reconcile avoid long-blocking calls?
Blocking workers on slow external APIs exhausts the controller-runtime worker pool. Long operations should set timeouts and leverage RequeueAfter for non-blocking polling.
Q1. How does DeltaFIFO guarantee event ordering in client-go?
DeltaFIFO maintains an ordered list of mutation deltas per object key. Rapid updates append to the key's delta slice, ensuring consumers never observe an Update before its associated Add event.
Q2. How do you implement robust Cloud Drift Detection without exceeding AWS API rate limits?
Combine RequeueAfter with randomized jitter, filter DescribeInstances queries by tags or IDs, and cache AWS responses to avoid hitting AWS Describe API throttling limits across large CR fleets.
Q3. How do you safely migrate a CRD across API versions (v1alpha1 to v1)?
Deploy a Conversion Webhook that converts representations on the fly, set both versions as served while preserving one storage version, migrate stored data with a storage version migrator, and deprecate the older version.
Q4. What is the difference between Server-Side Apply and r.Update()?
Server-Side Apply tracks field managers, allowing multiple controllers to manage distinct fields of the same resource declaratively without clobbering each other's updates.
Q5. Why is queue.Done(key) mandatory in client-go workqueues?
Get() marks a key as "in-flight" to prevent concurrent workers from processing the same key. Forgetting Done() permanently blocks subsequent updates for that key from ever being delivered.
Q6. How does controller-runtime prevent thundering herds on startup?
It blocks worker execution until cache.WaitForCacheSync verifies that all Informers have completed their initial LIST operation, preventing reconciles against half-populated caches.
Q7. How do you handle partial failures during external cloud resource creation?
Persist the external resource ID in .status or tags as early as possible. If a subsequent step fails, the next reconcile reads the persisted ID and adopts the existing cloud resource rather than re-creating it.
Q8. How does client-go workqueue rate limiting combine Token Bucket and Exponential Backoff?
It applies per-item exponential backoff (e.g. 5ms to 1000s) alongside a global token-bucket limiter (e.g. 10 QPS, 100 burst) to protect both the worker process and the target API server from saturation.
Q9. What are Mutating and Validating Admission Webhooks?
HTTP callbacks invoked by the API server prior to etcd persistence. Mutating webhooks inject defaults; validating webhooks enforce business rules and can reject invalid requests with descriptive errors.
Q10. How do you architect an operator to manage tens of thousands of custom resources?
Tune controller concurrency (MaxConcurrentReconciles), narrow Informer cache selectors via field or label selectors to avoid memory bloat, configure sharding, and use Server-Side Apply.
Q11. What is the risk of calling r.Client.Get with an uncached reader?
Bypassing the Informer cache via mgr.GetAPIReader() sends requests directly to the API server and etcd. Doing this inside high-frequency reconcile loops can overwhelm the Kubernetes control plane.
Q12. How do you guarantee zero-downtime upgrades of an operator controller?
Run multiple replicas with Leader Election enabled, set a RollingUpdate deployment strategy with maxUnavailable: 0, and ensure the reconcile loop is backwards-compatible with in-flight custom resources.
Q13. How do you avoid finalizer deadlocks during namespace deletion?
Ensure the controller continues running while namespaces are terminating. If the controller pod is deleted first, resources stuck with finalizers will block namespace deletion forever.
Q14. What are status Conditions and why follow the Metav1 Condition standard?
Conditions (type, status, reason, message, lastTransitionTime) provide a standardized, programmatic way for external tools (like ArgoCD, Helm, or CLI scripts) to determine resource health.
Q15. How does Crossplane or AWS Controllers for Kubernetes (ACK) scale across entire clouds?
They generate dedicated controllers per AWS service, utilize code-generation pipelines from cloud API specs, leverage AWS IAM roles per namespace, and maintain separate Informer caches per CRD group.