027 Project 27: Kubernetes Event Watcher
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.27Full 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
- Client from kubeconfig.
Events(ns).Watchwith ListOptions.- Range
ResultChan; type-assert*corev1.Event. - On channel close / error, backoff and reconnect.
- 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=0Stretch goals
- Filter
type=Warningonly. - JSON lines for log aggregation.
- Informer/factory instead of raw Watch.
- 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