-
Notifications
You must be signed in to change notification settings - Fork 20
/
service.go
116 lines (96 loc) · 3 KB
/
service.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
package continuousscanning
import (
"context"
"time"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
watch "k8s.io/apimachinery/pkg/watch"
"k8s.io/client-go/dynamic"
armoapi "github.com/armosec/armoapi-go/apis"
"github.com/kubescape/go-logger"
"github.com/kubescape/go-logger/helpers"
)
type ContinuousScanningService struct {
tl TargetLoader
shutdownRequested chan struct{}
workDone chan struct{}
k8sdynamic dynamic.Interface
eventHandlers []EventHandler
eventQueue *cooldownQueue
}
func (s *ContinuousScanningService) listen(ctx context.Context) <-chan armoapi.Command {
producedCommands := make(chan armoapi.Command)
listOpts := metav1.ListOptions{}
resourceEventsCh := make(chan watch.Event, 100)
gvrs := s.tl.LoadGVRs(ctx)
logger.L().Ctx(ctx).Info("fetched gvrs", helpers.Interface("gvrs", gvrs))
wp, _ := NewWatchPool(ctx, s.k8sdynamic, gvrs, listOpts)
wp.Run(ctx, resourceEventsCh)
logger.L().Ctx(ctx).Info("ran watch pool")
go func(shutdownCh <-chan struct{}, resourceEventsCh <-chan watch.Event, out *cooldownQueue) {
defer out.Stop(ctx)
for {
select {
case e := <-resourceEventsCh:
logger.L().Ctx(ctx).Debug(
"got event from channel",
helpers.Interface("event", e),
)
out.Enqueue(ctx, e)
case <-shutdownCh:
return
}
}
}(s.shutdownRequested, resourceEventsCh, s.eventQueue)
return producedCommands
}
func (s *ContinuousScanningService) work(ctx context.Context) {
for e := range s.eventQueue.ResultChan {
logger.L().Ctx(ctx).Debug(
"got an event to process",
helpers.Interface("event", e),
)
for idx := range s.eventHandlers {
handler := s.eventHandlers[idx]
err := handler.Handle(ctx, e)
if err != nil {
logger.L().Ctx(ctx).Error(
"failed to handle event",
helpers.Interface("event", e),
helpers.Error(err),
)
}
}
}
close(s.workDone)
}
// Launch launches the service.
//
// It sets up the provided watches, listens for events they deliver in the
// background and dispatches them to registered event handlers.
// Launch blocks until all the underlying watches are ready to accept events.
func (s *ContinuousScanningService) Launch(ctx context.Context) <-chan armoapi.Command {
out := make(chan armoapi.Command)
s.listen(ctx)
go s.work(ctx)
return out
}
func (s *ContinuousScanningService) AddEventHandler(fn EventHandler) {
s.eventHandlers = append(s.eventHandlers, fn)
}
func (s *ContinuousScanningService) Stop() {
close(s.shutdownRequested)
<-s.workDone
}
func NewContinuousScanningService(client dynamic.Interface, tl TargetLoader, queueSize int, sameEventCooldown time.Duration, h ...EventHandler) *ContinuousScanningService {
doneCh := make(chan struct{})
eventQueue := NewCooldownQueue(queueSize, sameEventCooldown)
workDone := make(chan struct{})
return &ContinuousScanningService{
tl: tl,
k8sdynamic: client,
shutdownRequested: doneCh,
eventHandlers: h,
eventQueue: eventQueue,
workDone: workDone,
}
}