mirror of
https://github.com/kubeshark/kubeshark.git
synced 2026-08-25 00:47:27 +00:00
* 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>
110 lines
2.3 KiB
Go
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
|
|
}
|
|
}
|
|
}
|