From 77c2d39de9bddcdc21dfce2ba4eb215f2988a67e Mon Sep 17 00:00:00 2001 From: Brandon Stoll Date: Thu, 9 Jul 2026 18:36:26 +0000 Subject: [PATCH] refactor: modernize Go API usage and error handling in core utilities (3/5) This is part 3/5 of an overall cleanup effort to fix linter issues and format files across the repository. In this step: - Modernize Go API usage and error handling across core utilities, events, exec, and pods. - Address linter issues in metrics and x/ packages. --- api/metallb/clientset/v1beta1/client.go | 10 ++++-- events/events.go | 6 ++-- events/status.go | 4 +-- events/status_test.go | 4 +-- exec/fake/fake_test.go | 10 +++--- exec/run/run_test.go | 8 ++--- flags/flags.go | 4 ++- logshim/logshim_test.go | 10 ++++-- metrics/metrics_test.go | 33 ++++++++++--------- pods/pods.go | 16 ++++----- pods/status.go | 4 +-- pods/status_test.go | 2 +- .../addcontainer/addcontainer_test.go | 1 - x/webhook/main.go | 4 ++- x/wire/file/client/main.go | 6 ++-- x/wire/forward/main.go | 6 ++-- x/wire/intf/client/main.go | 6 ++-- x/wire/wire.go | 6 +++- 18 files changed, 82 insertions(+), 58 deletions(-) diff --git a/api/metallb/clientset/v1beta1/client.go b/api/metallb/clientset/v1beta1/client.go index 0c88fbb6c..1287c69bb 100644 --- a/api/metallb/clientset/v1beta1/client.go +++ b/api/metallb/clientset/v1beta1/client.go @@ -92,10 +92,14 @@ var ( ) func init() { - metallbv1.AddToScheme(Scheme) + if err := metallbv1.AddToScheme(Scheme); err != nil { + panic(err) + } metav1.AddToGroupVersion(Scheme, groupVersion) - metav1.AddMetaToScheme(Scheme) + if err := metav1.AddMetaToScheme(Scheme); err != nil { + panic(err) + } } func GV() *schema.GroupVersion { @@ -105,7 +109,7 @@ func GV() *schema.GroupVersion { // NewForConfig returns a new Clientset based on c. func NewForConfig(c *rest.Config) (*Clientset, error) { config := *c - config.ContentConfig.GroupVersion = &groupVersion + config.GroupVersion = &groupVersion config.APIPath = "/apis" config.NegotiatedSerializer = scheme.Codecs.WithoutConversion() config.UserAgent = rest.DefaultKubernetesUserAgent() diff --git a/events/events.go b/events/events.go index 15ee0cec2..62218a338 100644 --- a/events/events.go +++ b/events/events.go @@ -102,9 +102,9 @@ func WatchEventStatus(ctx context.Context, client kubernetes.Interface, namespac // EventToStatus returns a pointer to a new EventStatus for an event. func EventToStatus(event *corev1.Event) *EventStatus { s := EventStatus{ - Name: event.ObjectMeta.Name, - Namespace: event.ObjectMeta.Namespace, - UID: event.ObjectMeta.UID, + Name: event.Name, + Namespace: event.Namespace, + UID: event.UID, } event.DeepCopyInto(&s.Event) event = &s.Event diff --git a/events/status.go b/events/status.go index 2df348290..1c9a3168a 100644 --- a/events/status.go +++ b/events/status.go @@ -63,7 +63,7 @@ func newWatcher(ctx context.Context, cancel func(), ch chan *EventStatus, stop f return w } -// SetProgress determins if progress output should be displayed while watching. +// SetProgress determines if progress output should be displayed while watching. func (w *Watcher) SetProgress(value bool) { w.mu.Lock() w.progress = value @@ -126,7 +126,7 @@ func (w *Watcher) isEventNormal(s *EventStatus) bool { for _, m := range errorMsgs { // Error out if message contains predefined message if strings.Contains(message, m) { - w.errCh <- fmt.Errorf("Event failed due to %s . Message: %s", s.Event.Reason, message) + w.errCh <- fmt.Errorf("event failed due to %s. message: %s", s.Event.Reason, message) w.cancel() return false } diff --git a/events/status_test.go b/events/status_test.go index 9dbdce136..bb70c741f 100644 --- a/events/status_test.go +++ b/events/status_test.go @@ -266,7 +266,7 @@ func TestIsEventNormal(t *testing.T) { want: ` 01:23:45 NS: ns1 Event name: event2 Type: Warning Message: 0/1 nodes are available: 1 Insufficient cpu. preemption: 0/1 nodes are available: 1 No preemption victims found for incoming pod.. `[1:], - errch: "Event failed due to . Message: 0/1 nodes are available: 1 Insufficient cpu. preemption: 0/1 nodes are available: 1 No preemption victims found for incoming pod..", + errch: "event failed due to . message: 0/1 nodes are available: 1 Insufficient cpu. preemption: 0/1 nodes are available: 1 No preemption victims found for incoming pod..", canceled: true, }, { @@ -275,7 +275,7 @@ func TestIsEventNormal(t *testing.T) { want: ` 01:23:45 NS: ns1 Event name: event3 Type: Warning Message: 0/1 nodes are available: 1 Insufficient memory. preemption: 0/1 nodes are available: 1 No preemption victims found for incoming pod.. `[1:], - errch: "Event failed due to . Message: 0/1 nodes are available: 1 Insufficient memory. preemption: 0/1 nodes are available: 1 No preemption victims found for incoming pod..", + errch: "event failed due to . message: 0/1 nodes are available: 1 Insufficient memory. preemption: 0/1 nodes are available: 1 No preemption victims found for incoming pod..", canceled: true, }, } { diff --git a/exec/fake/fake_test.go b/exec/fake/fake_test.go index 351518788..3b60c125b 100644 --- a/exec/fake/fake_test.go +++ b/exec/fake/fake_test.go @@ -284,11 +284,11 @@ func TestFailed(t *testing.T) { } } for _, u := range cmds.unexpected { - var ue Response + var unexpectedResp Response if len(tt.unexpected) > 0 { - ue = tt.unexpected[0] + unexpectedResp = tt.unexpected[0] } - t.Logf("Compare %v and %v", u, ue) + t.Logf("Compare %v and %v", u, unexpectedResp) if len(tt.unexpected) > 0 && tt.unexpected[0].String() == u.String() { tt.unexpected = tt.unexpected[1:] continue @@ -385,7 +385,9 @@ unexpected executions: c := cmds.Command(cmd.Cmd, cmd.Args...) c.SetStdout(&stdout) c.SetStderr(&stderr) - _ = c.Run() + if err := c.Run(); err != nil { + t.Errorf("c.Run() returned unexpected error: %v", err) + } } err := cmds.Done() if err.Error() != tt.done { diff --git a/exec/run/run_test.go b/exec/run/run_test.go index b57a6bc92..71c784571 100644 --- a/exec/run/run_test.go +++ b/exec/run/run_test.go @@ -111,11 +111,11 @@ func TestRunCommand(t *testing.T) { if string(got) != tt.want { t.Errorf("runCommand() got output %v, want %v", string(got), tt.want) } - if string(infos.Bytes()) != tt.wantInfos { - t.Errorf("runCommand() got info logs %v, want %v", string(infos.Bytes()), tt.wantInfos) + if infos.String() != tt.wantInfos { + t.Errorf("runCommand() got info logs %v, want %v", infos.String(), tt.wantInfos) } - if string(warnings.Bytes()) != tt.wantWarnings { - t.Errorf("runCommand() got warning logs %v, want %v", string(warnings.Bytes()), tt.wantWarnings) + if warnings.String() != tt.wantWarnings { + t.Errorf("runCommand() got warning logs %v, want %v", warnings.String(), tt.wantWarnings) } }) } diff --git a/flags/flags.go b/flags/flags.go index 765b9e354..6e58d3307 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -35,7 +35,9 @@ func Import(defmap map[string]string) { flag.Set("stderrthreshold", "INFO") //nolint:errcheck for k, v := range defmap { if f := flag.Lookup(k); f != nil { - f.Value.Set(v) + if err := f.Value.Set(v); err != nil { + klog.Warningf("Failed to set default value %q for flag %q: %v", v, k, err) + } f.DefValue = v } } diff --git a/logshim/logshim_test.go b/logshim/logshim_test.go index d63735c2a..114254235 100644 --- a/logshim/logshim_test.go +++ b/logshim/logshim_test.go @@ -36,10 +36,12 @@ Prefix:Line 3 want3 := want2 + `Prefix:partial line 2 ` // Writing with a partial line. - s.Write([]byte(`Line 1 + if _, err := s.Write([]byte(`Line 1 Line 2 Line 3 -partial line 1`)) +partial line 1`)); err != nil { + t.Fatalf("Write failed: %v", err) + } if got := b.String(); got != want1 { t.Errorf("First: Got %q, want %q", got, want1) } @@ -54,7 +56,9 @@ partial line 1`)) } // Write another partial line and close the shim - s.Write([]byte(`partial line 2`)) + if _, err := s.Write([]byte(`partial line 2`)); err != nil { + t.Fatalf("Write failed: %v", err) + } s.Close() if got := b.String(); got != want3 { t.Errorf("Second: Got %q, want %q", got, want3) diff --git a/metrics/metrics_test.go b/metrics/metrics_test.go index 5a8501ec7..f31c4db35 100644 --- a/metrics/metrics_test.go +++ b/metrics/metrics_test.go @@ -25,6 +25,7 @@ import ( epb "github.com/openconfig/kne/proto/event" "google.golang.org/api/option" "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" "google.golang.org/protobuf/proto" ) @@ -36,7 +37,7 @@ func newTestReporter(t *testing.T, ctx context.Context) (*Reporter, *pstest.Serv // Start a fake server running locally. srv := pstest.NewServer() // Connect to the server without using TLS. - conn, err := grpc.Dial(srv.Addr, grpc.WithInsecure()) + conn, err := grpc.NewClient(srv.Addr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { srv.Close() t.Fatalf("failed to start fake PubSub server: %v", err) @@ -93,9 +94,9 @@ func TestReportDeployClusterStart(t *testing.T) { if _, ok := e.Event.(*epb.KNEEvent_DeployClusterStart); !ok { t.Fatalf("event is not a DeployClusterStart event") } - te := e.GetDeployClusterStart() - if te.Cluster.Cluster != epb.Cluster_CLUSTER_TYPE_EXTERNAL { - t.Errorf("event has wrong cluster type, got %v, want %v", te.Cluster.Cluster, epb.Cluster_CLUSTER_TYPE_EXTERNAL) + gotEvent := e.GetDeployClusterStart() + if gotEvent.Cluster.Cluster != epb.Cluster_CLUSTER_TYPE_EXTERNAL { + t.Errorf("event has wrong cluster type, got %v, want %v", gotEvent.Cluster.Cluster, epb.Cluster_CLUSTER_TYPE_EXTERNAL) } } @@ -129,9 +130,9 @@ func TestReportDeployClusterEnd(t *testing.T) { if _, ok := e.Event.(*epb.KNEEvent_DeployClusterEnd); !ok { t.Fatalf("event is not a DeployClusterEnd event") } - te := e.GetDeployClusterEnd() - if te.Error != eventErr.Error() { - t.Errorf("event has wrong error message, got %v, want %v", te.Error, eventErr.Error()) + gotEvent := e.GetDeployClusterEnd() + if gotEvent.Error != eventErr.Error() { + t.Errorf("event has wrong error message, got %v, want %v", gotEvent.Error, eventErr.Error()) } } @@ -164,9 +165,9 @@ func TestReportCreateTopologyStart(t *testing.T) { if _, ok := e.Event.(*epb.KNEEvent_CreateTopologyStart); !ok { t.Fatalf("event is not a CreateTopologyStart event") } - te := e.GetCreateTopologyStart() - if te.Topology.LinkCount != 3 { - t.Errorf("event has wrong link count, got %v, want 3", te.Topology.LinkCount) + gotEvent := e.GetCreateTopologyStart() + if gotEvent.Topology.LinkCount != 3 { + t.Errorf("event has wrong link count, got %v, want 3", gotEvent.Topology.LinkCount) } } @@ -200,9 +201,9 @@ func TestReportCreateTopologyEnd(t *testing.T) { if _, ok := e.Event.(*epb.KNEEvent_CreateTopologyEnd); !ok { t.Fatalf("event is not a CreateTopologyEnd event") } - te := e.GetCreateTopologyEnd() - if te.Error != eventErr.Error() { - t.Errorf("event has wrong error message, got %v, want %v", te.Error, eventErr.Error()) + gotEvent := e.GetCreateTopologyEnd() + if gotEvent.Error != eventErr.Error() { + t.Errorf("event has wrong error message, got %v, want %v", gotEvent.Error, eventErr.Error()) } } @@ -212,7 +213,7 @@ func TestNewReporter(t *testing.T) { defer srv.Close() // Test missing topic - conn1, err := grpc.Dial(srv.Addr, grpc.WithInsecure()) + conn1, err := grpc.NewClient(srv.Addr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { t.Fatalf("failed to dial fake PubSub server: %v", err) } @@ -222,7 +223,7 @@ func TestNewReporter(t *testing.T) { } // Create default topic for default project/topic test - conn2, err := grpc.Dial(srv.Addr, grpc.WithInsecure()) + conn2, err := grpc.NewClient(srv.Addr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { t.Fatalf("failed to dial fake PubSub server: %v", err) } @@ -241,7 +242,7 @@ func TestNewReporter(t *testing.T) { client.Close() // Test NewReporter success with default project/topic - conn3, err := grpc.Dial(srv.Addr, grpc.WithInsecure()) + conn3, err := grpc.NewClient(srv.Addr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { t.Fatalf("failed to dial fake PubSub server: %v", err) } diff --git a/pods/pods.go b/pods/pods.go index a1d33b5f7..9a79b0852 100644 --- a/pods/pods.go +++ b/pods/pods.go @@ -121,11 +121,11 @@ func (c *ContainerStatus) String() string { } func (c *ContainerStatus) Equal(oc *ContainerStatus) bool { - return !(c.Name != oc.Name || - c.Image != oc.Image || - c.Ready != oc.Ready || - c.Reason != oc.Reason || - c.Message != oc.Message) + return c.Name == oc.Name && + c.Image == oc.Image && + c.Ready == oc.Ready && + c.Reason == oc.Reason && + c.Message == oc.Message } // GetPodStatus returns the status of the pods found in the supplied namespace. @@ -183,9 +183,9 @@ func WatchPodStatus(ctx context.Context, client kubernetes.Interface, namespace // PodToStatus returns a pointer to a new PodStatus for pod. func PodToStatus(pod *corev1.Pod) *PodStatus { s := PodStatus{ - Name: pod.ObjectMeta.Name, - Namespace: pod.ObjectMeta.Namespace, - UID: pod.ObjectMeta.UID, + Name: pod.Name, + Namespace: pod.Namespace, + UID: pod.UID, Phase: pod.Status.Phase, // Ready will be set to false below if one of the containers is not ready Ready: len(pod.Status.ContainerStatuses)+len(pod.Status.InitContainerStatuses) > 0, diff --git a/pods/status.go b/pods/status.go index 6f30310b7..73ac4bf8f 100644 --- a/pods/status.go +++ b/pods/status.go @@ -71,7 +71,7 @@ func newWatcher(ctx context.Context, cancel func(), ch chan *PodStatus, stop fun return w } -// SetProgress determins if progress output should be displayed while watching. +// SetProgress determines if progress output should be displayed while watching. func (w *Watcher) SetProgress(value bool) { w.mu.Lock() w.progress = value @@ -176,7 +176,7 @@ func (w *Watcher) updatePod(s *PodStatus) bool { w.podStates[s.UID] = newState } if newState == "failed" { - w.errCh <- fmt.Errorf("Pod %s failed to deploy", s.Name) + w.errCh <- fmt.Errorf("pod %s failed to deploy", s.Name) w.cancel() return false } diff --git a/pods/status_test.go b/pods/status_test.go index f127b4cd6..c38f2b611 100644 --- a/pods/status_test.go +++ b/pods/status_test.go @@ -100,7 +100,7 @@ func TestUpdatePod(t *testing.T) { want: ` 01:23:45 POD: pod1 is now failed `[1:], - errch: `Pod pod1 failed to deploy`, + errch: `pod pod1 failed to deploy`, canceled: true, }, { diff --git a/x/webhook/examples/addcontainer/addcontainer_test.go b/x/webhook/examples/addcontainer/addcontainer_test.go index e119fb1e8..892451ad8 100644 --- a/x/webhook/examples/addcontainer/addcontainer_test.go +++ b/x/webhook/examples/addcontainer/addcontainer_test.go @@ -69,5 +69,4 @@ func TestAddContainer(t *testing.T) { } }) } - } diff --git a/x/webhook/main.go b/x/webhook/main.go index 0df6bf624..9d50dab95 100644 --- a/x/webhook/main.go +++ b/x/webhook/main.go @@ -90,7 +90,9 @@ func parseRequest(r http.Request) (*admissionv1.AdmissionReview, error) { } bodybuf := new(bytes.Buffer) - bodybuf.ReadFrom(r.Body) + if _, err := bodybuf.ReadFrom(r.Body); err != nil { + return nil, fmt.Errorf("failed to read request body: %w", err) + } body := bodybuf.Bytes() if len(body) == 0 { diff --git a/x/wire/file/client/main.go b/x/wire/file/client/main.go index 4f9d56f35..4a5cc0a9b 100644 --- a/x/wire/file/client/main.go +++ b/x/wire/file/client/main.go @@ -48,7 +48,7 @@ func main() { opts := []grpc.DialOption{ grpc.WithTransportCredentials(insecure.NewCredentials()), } - conn, err := grpc.DialContext(ctx, *addr, opts...) + conn, err := grpc.NewClient(*addr, opts...) if err != nil { log.Fatalf("Failed to dial %q: %v", *addr, err) } @@ -64,7 +64,9 @@ func main() { return err } defer func() { - stream.CloseSend() + if err := stream.CloseSend(); err != nil { + log.Warningf("Failed to close send stream: %v", err) + } }() log.Infof("Transmitting endpoint %v over wire...", e) return w.Transmit(ctx, stream) diff --git a/x/wire/forward/main.go b/x/wire/forward/main.go index 15d6eb7e8..5d5bdf8ce 100644 --- a/x/wire/forward/main.go +++ b/x/wire/forward/main.go @@ -122,7 +122,7 @@ func main() { opts := []grpc.DialOption{ grpc.WithTransportCredentials(insecure.NewCredentials()), } - conn, err := grpc.DialContext(ctx, a, opts...) + conn, err := grpc.NewClient(a, opts...) if err != nil { log.Fatalf("Failed to dial %q: %v", a, err) } @@ -138,7 +138,9 @@ func main() { return err } defer func() { - stream.CloseSend() + if err := stream.CloseSend(); err != nil { + log.Warningf("Failed to close send stream: %v", err) + } }() log.Infof("Transmitting endpoint %v over wire...", e) return w.Transmit(ctx, stream) diff --git a/x/wire/intf/client/main.go b/x/wire/intf/client/main.go index 1ce877551..68341c882 100644 --- a/x/wire/intf/client/main.go +++ b/x/wire/intf/client/main.go @@ -48,7 +48,7 @@ func main() { opts := []grpc.DialOption{ grpc.WithTransportCredentials(insecure.NewCredentials()), } - conn, err := grpc.DialContext(ctx, *addr, opts...) + conn, err := grpc.NewClient(*addr, opts...) if err != nil { log.Fatalf("Failed to dial %q: %v", *addr, err) } @@ -64,7 +64,9 @@ func main() { return err } defer func() { - stream.CloseSend() + if err := stream.CloseSend(); err != nil { + log.Warningf("Failed to close send stream: %v", err) + } }() log.Infof("Transmitting endpoint %v over wire...", e) return w.Transmit(ctx, stream) diff --git a/x/wire/wire.go b/x/wire/wire.go index 219e62659..30dacb1e7 100644 --- a/x/wire/wire.go +++ b/x/wire/wire.go @@ -151,7 +151,11 @@ func (w *Wire) Transmit(ctx context.Context, stream Stream) error { }) g.Go(func() error { if cs, ok := stream.(grpc.ClientStream); ok { - defer cs.CloseSend() + defer func() { + if err := cs.CloseSend(); err != nil { + log.Warningf("Failed to close send stream: %v", err) + } + }() } for { data, err := w.src.Read()