027 Project 27: Kubernetes Event Watcher

Updated

September 8, 2026

027 Build a Kubernetes Event Watcher

Stream cluster Events in real time with client-go Watch APIs. Print concise, timestamped lines operators can follow during incident response or CI rollout debugging.

kubeconfig -> clientset -> Events.Watch -> ResultChan -> print

Problem statement

CLI kevents:

  • Watch Events in one namespace or all
  • Optional field selector (e.g. involved object kind)
  • Print: TIME TYPE REASON OBJECT MESSAGE
  • Handle watch reconnects on error
  • Graceful shutdown on SIGINT

Acceptance criteria

  • Opens a Watch and prints events until interrupted
  • Namespace flag works
  • Context cancel stops cleanly (Stop() watch)
  • Errors are visible; process does not silent-exit on first blip without log
  • Builds with modules

Setup

mkdir kevents && cd kevents
go mod init example.com/kevents
go get k8s.io/client-go@latest
go get k8s.io/apimachinery@latest
go mod tidy
# go 1.27

Full main.go

package main

import (
    "context"
    "flag"
    "fmt"
    "os"
    "os/signal"
    "path/filepath"
    "syscall"
    "time"

    corev1 "k8s.io/api/core/v1"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/apimachinery/pkg/watch"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/clientcmd"
    "k8s.io/client-go/util/homedir"
)

func formatEvent(e *corev1.Event) string {
    obj := e.InvolvedObject.Kind + "/" + e.InvolvedObject.Name
    if e.InvolvedObject.Namespace != "" {
        obj = e.InvolvedObject.Namespace + "/" + obj
    }
    ts := e.LastTimestamp.Time
    if ts.IsZero() {
        ts = e.EventTime.Time
    }
    if ts.IsZero() {
        ts = time.Now()
    }
    return fmt.Sprintf("%s %-12s %-20s %-40s %s",
        ts.Format(time.RFC3339),
        e.Type,
        e.Reason,
        obj,
        e.Message,
    )
}

func main() {
    var kubeconfig string
    if home := homedir.HomeDir(); home != "" {
        kubeconfig = filepath.Join(home, ".kube", "config")
    }
    kc := flag.String("kubeconfig", kubeconfig, "kubeconfig path")
    ns := flag.String("n", "", "namespace (empty=all)")
    fieldSelector := flag.String("field-selector", "", "field selector")
    flag.Parse()

    cfg, err := clientcmd.BuildConfigFromFlags("", *kc)
    if err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }
    cli, err := kubernetes.NewForConfig(cfg)
    if err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }

    ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
    defer stop()

    namespace := *ns
    if namespace == "" {
        namespace = metav1.NamespaceAll
    }

    fmt.Fprintln(os.Stderr, "watching events; Ctrl+C to stop")

    for {
        if ctx.Err() != nil {
            return
        }
        w, err := cli.CoreV1().Events(namespace).Watch(ctx, metav1.ListOptions{
            FieldSelector: *fieldSelector,
        })
        if err != nil {
            fmt.Fprintln(os.Stderr, "watch error:", err)
            select {
            case <-ctx.Done():
                return
            case <-time.After(2 * time.Second):
            }
            continue
        }

        func() {
            defer w.Stop()
            for {
                select {
                case <-ctx.Done():
                    return
                case ev, ok := <-w.ResultChan():
                    if !ok {
                        fmt.Fprintln(os.Stderr, "watch closed; reconnecting")
                        return
                    }
                    if ev.Type == watch.Error {
                        fmt.Fprintf(os.Stderr, "watch event error: %v\n", ev.Object)
                        return
                    }
                    e, ok := ev.Object.(*corev1.Event)
                    if !ok {
                        fmt.Printf("%s %T\n", ev.Type, ev.Object)
                        continue
                    }
                    fmt.Printf("%-8s %s\n", ev.Type, formatEvent(e))
                }
            }
        }()
    }
}

Step-by-step build path

  1. Client from kubeconfig.
  2. Events(ns).Watch with ListOptions.
  3. Range ResultChan; type-assert *corev1.Event.
  4. On channel close / error, backoff and reconnect.
  5. SIGINT via context cancels Watch.

Run and verification

go run .
go run . -n default
# generate events:
kubectl run wget --image=busybox --restart=Never -- wget -O- https://kubernetes.default
kubectl delete pod wget --force --grace-period=0

Stretch goals

  1. Filter type=Warning only.
  2. JSON lines for log aggregation.
  3. Informer/factory instead of raw Watch.
  4. Metrics: events/sec counter.

Pitfalls

Pitfall Fix
Not calling w.Stop() leak; always defer Stop
Ignoring watch close reconnect loop
Printing only %T format useful fields
No cancel hang on Ctrl+C

Learning goals

  • Watch API lifecycle
  • Operator-facing event streams
  • Reconnect patterns for long-lived clients