Files
kubeshark/kubernetes/watch.go
T
Volodymyr StoikoandAlon Girmonsky f38980ea94 deps: bump indirect deps to clear critical/high Dependabot alerts (#1952)
* deps: bump indirect deps to clear critical/high Dependabot alerts

Bumps the vulnerable indirect dependencies flagged as critical or high
severity in Dependabot:

- golang.org/x/crypto  v0.39.0 -> v0.54.0  (7 critical + 2 high: SSH agent
  constraint/key-constraint bypass, @revoked auth bypass, FIDO/U2F presence
  check bypass, VerifiedPublicKeyCallback permission skip, infinite loop on
  large channel writes, client-induced server deadlock, RSA/DSA DoS, byte
  arithmetic underflow panic)
- google.golang.org/grpc  v1.68.1 -> v1.83.0  (critical: authz bypass via
  missing leading slash in :path; high: xDS RBAC and HTTP/2 issues)
- github.com/containerd/containerd  v1.7.27 -> v1.7.34  (high: LABEL ->
  restart-monitor binary:// host-root RCE, runAsNonRoot evasion, local
  privesc via wide CRI directory permissions)
- oras.land/oras-go/v2  v2.6.0 -> v2.6.2  (high: CVE-2026-50163 hardlink
  extract-dir escape, credential forwarding via unvalidated Location header)
- github.com/moby/spdystream  v0.5.0 -> v0.5.1  (high: DoS on CRI)

Transitively pulls up x/net, x/sync, x/sys, x/term, x/text, x/time,
x/oauth2, protobuf, filepath-securejoin, selinux and go-logr via go mod tidy.

The go directive moves 1.24.0 -> 1.25.0 (required by the upgraded modules);
the explicit toolchain pin is dropped. CI resolves Go from go.mod, so no
workflow changes are needed.

go build ./... and go test ./... pass.

* ci: move golangci-lint to v2, fix resulting lint issues

golangci-lint-action@v3 pins `latest` to v1.64.8, which is built with
go1.24 and refuses to run now that go.mod targets 1.25.0:

  can't load config: the Go language version (go1.24) used to build
  golangci-lint is lower than the targeted Go version (1.25.0)

Move the job to golangci-lint-action@v7 + v2.8.0 and add a .golangci.yml
mirroring the hub repo's v2 config: govet, staticcheck, ineffassign and
unused, plus gofmt/goimports as formatters.

Fixes for the issues that surfaced:
- ST1005: lowercase error strings, drop trailing '!' in connect/hub.go
- SA4011: kubernetes/watch.go had a `break` inside a `select` default
  that broke the select rather than the loop, i.e. a no-op; removed
- QF1008: drop the embedded ChartPathOptions selector in helm.go
- QF1003: tagged switch on r.URL.Path in mcp_test.go
- QF1004: strings.Replace(..., -1) -> strings.ReplaceAll
- gofmt -s and goimports with a local prefix across the tree

errcheck is not in the enabled set, matching hub.

* cmd: clarify --time parse error in pcap dump

The error neither named the offending flag/value nor separated the
wrapped error from the message. Reported by Copilot on #1952.

---------

Co-authored-by: Alon Girmonsky <1990761+alongir@users.noreply.github.com>
2026-08-04 08:35:41 -07:00

110 lines
2.3 KiB
Go

package kubernetes
import (
"context"
"errors"
"fmt"
"sync"
"time"
"github.com/rs/zerolog/log"
"k8s.io/apimachinery/pkg/watch"
"github.com/kubeshark/kubeshark/debounce"
)
type EventFilterer interface {
Filter(*WatchEvent) (bool, error)
}
type WatchCreator interface {
NewWatcher(ctx context.Context, namespace string) (watch.Interface, error)
}
func FilteredWatch(ctx context.Context, watcherCreator WatchCreator, targetNamespaces []string, filterer EventFilterer) (<-chan *WatchEvent, <-chan error) {
eventChan := make(chan *WatchEvent)
errorChan := make(chan error)
var wg sync.WaitGroup
for _, targetNamespace := range targetNamespaces {
wg.Add(1)
go func(targetNamespace string) {
defer wg.Done()
watchRestartDebouncer := debounce.NewDebouncer(1*time.Minute, func() {})
for {
watcher, err := watcherCreator.NewWatcher(ctx, targetNamespace)
if err != nil {
errorChan <- fmt.Errorf("error in k8s watch: %v", err)
break
}
err = startWatchLoop(ctx, watcher, filterer, eventChan) // blocking
watcher.Stop()
select {
case <-ctx.Done():
return
default:
}
if err != nil {
errorChan <- fmt.Errorf("error in k8s watch: %v", err)
break
} else {
if !watchRestartDebouncer.IsOn() {
if err := watchRestartDebouncer.SetOn(); err != nil {
log.Error().Err(err).Send()
}
log.Warn().Msg("K8s watch channel closed, restarting watcher...")
time.Sleep(time.Second * 5)
continue
} else {
errorChan <- errors.New("K8s watch unstable, closes frequently")
break
}
}
}
}(targetNamespace)
}
go func() {
<-ctx.Done()
wg.Wait()
close(eventChan)
close(errorChan)
}()
return eventChan, errorChan
}
func startWatchLoop(ctx context.Context, watcher watch.Interface, filterer EventFilterer, eventChan chan<- *WatchEvent) error {
resultChan := watcher.ResultChan()
for {
select {
case e, isChannelOpen := <-resultChan:
if !isChannelOpen {
return nil
}
wEvent := WatchEvent(e)
if wEvent.Type == watch.Error {
return wEvent.ToError()
}
if pass, err := filterer.Filter(&wEvent); err != nil {
return err
} else if !pass {
continue
}
eventChan <- &wEvent
case <-ctx.Done():
return nil
}
}
}