Applied Go • Cloud Infrastructure • Platform Engineering

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.

What's an operator Project setup The EC2 CRD Admission Webhooks Interactive simulator The reconcile loop Reconcile timeline Finalizers & cloud cleanup AWS SDK v2 integration Controller wiring Testing with envtest Deploying to cluster Raw client-go internals Interview Questions
Verified against real tooling & live AWS APIs

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.

The golden rule: reconciliation is level-based, not edge-based

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.

Terminal
$ mkdir ec2-operator $ cd ec2-operator $ go mod init github.com/r4rajat/ec2-operator # 1. Initialize Kubebuilder scaffolding $ kubebuilder init --domain cloud.com --repo github.com/r4rajat/ec2-operator # 2. Scaffold the Ec2Instance CRD and Controller $ kubebuilder create api --group compute --version v1 --kind Ec2Instance Create Resource [y/n] y Create Controller [y/n] y # 3. Download the AWS SDK for Go v2 dependencies $ go get github.com/aws/aws-sdk-go-v2 $ go get github.com/aws/aws-sdk-go-v2/config $ go get github.com/aws/aws-sdk-go-v2/credentials $ go get github.com/aws/aws-sdk-go-v2/service/ec2

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.

api/v1/ec2instance_types.go
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"`
}
Why the /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 .status does not increment metadata.generation. This prevents the controller's own status updates from infinitely triggering new reconcile loops.
  • RBAC Isolation: Developers can be granted permissions to edit .spec without having write access to spoof .status, ensuring cluster integrity.

Running make manifests automatically compiles the Go types into the OpenAPI v3 YAML manifest:

config/crd/bases/compute.cloud.com_ec2instances.yaml (abridged)
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.

The Crucial Distinction: Webhooks vs. Reconcilers

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 kubectl to return. If a Validating Webhook rejects a request, nothing is written to etcd, the change never happened, and the Reconciler is never triggered.
  • Reconciliation is ASYNCHRONOUS and Post-Persistence: The Reconcile loop 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:

API Request Pipeline
kubectl apply / delete │ ▼ ┌───────────────────────────────┐ │ 1. Authentication & RBAC AuthZ │ ← Who are you? Do you have permissions? └──────────────┬────────────────┘ │ ▼ ┌───────────────────────────────┐ │ 2. Mutating Webhooks │ ← Can modify object (defaults, tag injection, sidecars) └──────────────┬────────────────┘ │ ▼ ┌───────────────────────────────┐ │ 3. OpenAPI v3 Schema Validation│ ← Type checks, required fields, regexes └──────────────┬────────────────┘ │ ▼ ┌───────────────────────────────┐ │ 4. Validating Webhooks │ ← Semantic checks (cross-field rules, immutability, guards) └──────────────┬────────────────┘ │ ├── (Rejected) ──► Returns HTTP 400/422/403 (Nothing saved) │ ▼ (Approved) ┌───────────────────────────────┐ │ 5. Persistence in etcd │ ← Object saved! API server returns HTTP 200/201 └──────────────┬────────────────┘ │ ▼ (Watch Event Emitted) ┌───────────────────────────────┐ │ 6. Go Reconciler Loop │ ← Asynchronous: Calls AWS SDK to provision/terminate VM └───────────────────────────────┘

Scaffolding Webhooks with Kubebuilder

Kubebuilder simplifies webhook creation by scaffolding the webhook interfaces, registration logic, and Kubernetes manifests with a single command:

Terminal
$ kubebuilder create webhook \ --group compute \ --version v1 \ --kind Ec2Instance \ --defaulting \ --programmatic-validation

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:

api/v1/ec2instance_webhook.go (Mutating Defaulter)
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:

api/v1/ec2instance_webhook.go (Validating Validator)
// +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:

kube-apiserver • admissionreview Admission Webhook Pipeline Simulator
Synchronous Pre-Persistence
Step 1
💻
Client Request
HTTP POST/DELETE
→
Step 2
🪄
Mutating Webhook
Mutated
→
Step 3
📋
OpenAPI Schema
Valid
→
Step 4
🛡️
Validating Webhook
Approved
→
Step 5
💾
etcd Storage
Saved (201)
✓ HTTP 201 Created: Mutating Webhook injected default region, instanceType, and org tags. Validated and persisted to etcd!
Incoming Client Payload
apiVersion: compute.cloud.com/v1
kind: Ec2Instance
metadata:
  name: web-server
spec:
  amiId: ami-0c55b159cbfafe1f0
Admission Result / Error Response
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
The Mutating Admission Webhook intercepted the incoming JSON payload, defaulted missing fields, and injected mandatory compliance tags before etcd persistence.

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:

config/webhook/manifests.yaml (with cert-manager CA injection)
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
Production Pitfall: 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 namespaceSelector on your webhook configuration to exclude kube-system and 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:

controller-runtime • AWS SDK v2 Cloud Operator Control Loop Simulator
Step 0 of 5: Cluster Ready
☸ Kubernetes API & etcd
Control Plane
Resource:—
spec.instanceType:—
spec.region:—
metadata.finalizers:[]
deletionTimestamp:null
status.instanceId:—
status.state:None
status.publicIP:—
⚙ Go Reconciler Loop
controller-runtime
Active Worker:Idle
WorkQueue:[] (0 items)
Waiting for watch event...
☁ AWS EC2 Cloud
us-east-1
Instance ID:—
AWS State:No Instance
Public IPv4:—
Hourly Cost:$0.00 / hr
Click “1. Apply EC2Instance CR” to submit a custom resource declaring a desired AWS EC2 instance in us-east-1.

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.

internal/controller/ec2instance_controller.go
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?

1
Reconcile #1: Creation & Waiter (Runs ~15s)
Triggered by CR creation. Reconcile #1 calls 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().
2
Reconcile #2: Idempotency in Action
Starts only after Reconcile #1 completes. When Reconcile #2 calls 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!
3
Reconcile #3: Workqueue Coalescing
Triggered by the earlier status update. If Reconcile #2 was already queued when the status update arrived, controller-runtime's workqueue de-duplicates identical keys. Even if processed, it observes unchanged state and immediately enters idle polling mode.

Try the single-field drift visualizer below to see level-based reconciliation in action:

Desired Spec
instanceTypet3.micro
regionus-east-1
↻
Observed AWS State
staterunning
Click "Simulate manual drift" to simulate an out-of-band AWS console change, then "Reconcile" to watch the operator correct it.

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:

Parent CR (In-Cluster)
ConfigMap
Pod
Click “Delete Parent Resource” to see what the Kubernetes garbage collector does to in-cluster children.
The fatal mistake: assuming Kubernetes garbage-collects AWS

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:

internal/controller/aws_helpers.go
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:

internal/controller/ec2instance_controller.go
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)
}
Understanding Requeue Return Patterns

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.

internal/controller/ec2instance_controller_test.go (excerpt)
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())
    })
  })
})
$ make test
Running Suite: Ec2Instance Controller Suite - /internal/controller =================================================================== Random Seed: 1726084496 Will run 5 of 5 specs ••••• Ran 5 of 5 Specs in 4.812 seconds SUCCESS! -- 5 Passed | 0 Failed | 0 Pending | 0 Skipped PASS ok github.com/r4rajat/ec2-operator/internal/controller 5.120s coverage: 88.4% of statements

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:

Deployment sequence
# 1. Create a local Kind test cluster $ kind create cluster --name operator-lab # 2. Build and load the operator container $ make docker-build IMG=ec2-operator:latest $ kind load docker-image ec2-operator:latest --name operator-lab # 3. Apply AWS credentials Secret (In EKS, use IAM Roles for Service Accounts instead) $ kubectl create secret generic aws-credentials \ --from-literal=AWS_ACCESS_KEY_ID=$AWS_ACCESS_KEY_ID \ --from-literal=AWS_SECRET_ACCESS_KEY=$AWS_SECRET_ACCESS_KEY # 4. Install CRDs and deploy the operator controller-manager $ make install $ make deploy IMG=ec2-operator:latest

Now, declare your AWS infrastructure using standard Kubernetes YAML:

config/samples/compute_v1_ec2instance.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
$ kubectl apply -f config/samples/compute_v1_ec2instance.yaml && kubectl get ec2instances -w
NAME INSTANCETYPE STATE PUBLICIP INSTANCEID AGE web-server-prod t3.micro pending i-078a9bc4e12fa89b1 8s web-server-prod t3.micro running 54.214.88.19 i-078a9bc4e12fa89b1 34s

Verify against the real AWS CLI:

$ aws ec2 describe-instances --instance-ids i-078a9bc4e12fa89b1 --query "Reservations[].Instances[].[InstanceId,State.Name,PublicIpAddress]" --output text
i-078a9bc4e12fa89b1 running 54.214.88.19

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:

Under the hood: client-go Informer and WorkQueue wiring
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.

← Back to the roadmap