feat(monitor/collectors/kev): 1h KEV collector with set-diff and chime topic

This commit is contained in:
CarterPerez-dev 2026-05-01 22:44:37 -04:00
parent 7395f1064f
commit bdca22ee50
2 changed files with 286 additions and 0 deletions

View File

@ -0,0 +1,140 @@
// ©AngelaMos | 2026
// collector.go
package kev
import (
"context"
"encoding/json"
"log/slog"
"time"
"github.com/carterperez-dev/monitor-the-situation/backend/internal/events"
)
const (
Name = "kev"
defaultKEVCadence = time.Hour
dateLayout = "2006-01-02"
)
type Fetcher interface {
FetchCatalog(ctx context.Context) (Catalog, error)
}
type Repository interface {
Insert(ctx context.Context, row Row) error
KnownIDs(ctx context.Context, ids []string) (map[string]bool, error)
}
type Emitter interface {
Emit(ev events.Event)
}
type StateRecorder interface {
RecordSuccess(ctx context.Context, name string, eventCount int64) error
RecordError(ctx context.Context, name, errMsg string) error
}
type CollectorConfig struct {
Interval time.Duration
Fetcher Fetcher
Repo Repository
Emitter Emitter
State StateRecorder
Logger *slog.Logger
}
type Collector struct {
cfg CollectorConfig
logger *slog.Logger
}
func NewCollector(cfg CollectorConfig) *Collector {
if cfg.Interval <= 0 {
cfg.Interval = defaultKEVCadence
}
if cfg.Logger == nil {
cfg.Logger = slog.Default()
}
return &Collector{cfg: cfg, logger: cfg.Logger}
}
func (c *Collector) Name() string { return Name }
func (c *Collector) Run(ctx context.Context) error {
ticker := time.NewTicker(c.cfg.Interval)
defer ticker.Stop()
c.tick(ctx)
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
c.tick(ctx)
}
}
}
func (c *Collector) tick(ctx context.Context) {
cat, err := c.cfg.Fetcher.FetchCatalog(ctx)
if err != nil {
c.logger.Warn("kev fetch", "err", err)
_ = c.cfg.State.RecordError(ctx, Name, err.Error())
return
}
ids := make([]string, 0, len(cat.Vulnerabilities))
for _, v := range cat.Vulnerabilities {
ids = append(ids, v.CveID)
}
known, err := c.cfg.Repo.KnownIDs(ctx, ids)
if err != nil {
c.logger.Warn("kev known ids", "err", err)
_ = c.cfg.State.RecordError(ctx, Name, err.Error())
return
}
now := time.Now().UTC()
emitted := int64(0)
for _, v := range cat.Vulnerabilities {
if known[v.CveID] {
continue
}
dateAdded, _ := time.Parse(dateLayout, v.DateAdded)
var dueDate *time.Time
if t, perr := time.Parse(dateLayout, v.DueDate); perr == nil {
dueDate = &t
}
raw, _ := json.Marshal(v)
row := Row{
CveID: v.CveID,
Vendor: v.VendorProject,
Product: v.Product,
VulnerabilityName: v.VulnerabilityName,
DateAdded: dateAdded,
DueDate: dueDate,
RansomwareUse: v.KnownRansomwareCampaignUse,
Payload: raw,
}
if ierr := c.cfg.Repo.Insert(ctx, row); ierr != nil {
c.logger.Warn("kev insert", "id", v.CveID, "err", ierr)
continue
}
c.cfg.Emitter.Emit(events.Event{
Topic: events.TopicKEVAdded,
Timestamp: now,
Source: Name,
Payload: json.RawMessage(raw),
})
emitted++
}
_ = c.cfg.State.RecordSuccess(ctx, Name, emitted)
}

View File

@ -0,0 +1,146 @@
// ©AngelaMos | 2026
// collector_test.go
package kev_test
import (
"context"
"sync"
"testing"
"time"
"github.com/stretchr/testify/require"
"github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/kev"
"github.com/carterperez-dev/monitor-the-situation/backend/internal/events"
)
type stubFetcher struct {
cat kev.Catalog
err error
}
func (s *stubFetcher) FetchCatalog(context.Context) (kev.Catalog, error) {
return s.cat, s.err
}
type stubKEVRepo struct {
mu sync.Mutex
inserted []string
known map[string]bool
}
func (r *stubKEVRepo) Insert(_ context.Context, row kev.Row) error {
r.mu.Lock()
defer r.mu.Unlock()
r.inserted = append(r.inserted, row.CveID)
if r.known == nil {
r.known = map[string]bool{}
}
r.known[row.CveID] = true
return nil
}
func (r *stubKEVRepo) KnownIDs(_ context.Context, ids []string) (map[string]bool, error) {
r.mu.Lock()
defer r.mu.Unlock()
out := make(map[string]bool)
for _, id := range ids {
if r.known[id] {
out[id] = true
}
}
return out, nil
}
func (r *stubKEVRepo) Inserted() []string {
r.mu.Lock()
defer r.mu.Unlock()
out := make([]string, len(r.inserted))
copy(out, r.inserted)
return out
}
type stubKEVEmitter struct {
mu sync.Mutex
events []events.Event
}
func (e *stubKEVEmitter) Emit(ev events.Event) {
e.mu.Lock()
defer e.mu.Unlock()
e.events = append(e.events, ev)
}
func (e *stubKEVEmitter) Events() []events.Event {
e.mu.Lock()
defer e.mu.Unlock()
out := make([]events.Event, len(e.events))
copy(out, e.events)
return out
}
type stubKEVState struct{}
func (stubKEVState) RecordSuccess(context.Context, string, int64) error { return nil }
func (stubKEVState) RecordError(context.Context, string, string) error { return nil }
func TestCollector_OnlyEmitsNewKEVs(t *testing.T) {
ftch := &stubFetcher{cat: kev.Catalog{Vulnerabilities: []kev.Vulnerability{
{CveID: "CVE-2024-OLD", VendorProject: "X", DateAdded: "2026-04-01"},
{CveID: "CVE-2024-NEW", VendorProject: "Y", DateAdded: "2026-05-01"},
}}}
repo := &stubKEVRepo{known: map[string]bool{"CVE-2024-OLD": true}}
emt := &stubKEVEmitter{}
c := kev.NewCollector(kev.CollectorConfig{
Interval: 30 * time.Millisecond,
Fetcher: ftch,
Repo: repo,
Emitter: emt,
State: stubKEVState{},
})
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
defer cancel()
_ = c.Run(ctx)
inserted := repo.Inserted()
require.NotEmpty(t, inserted)
for _, id := range inserted {
require.NotEqual(t, "CVE-2024-OLD", id)
require.Equal(t, "CVE-2024-NEW", id)
}
evs := emt.Events()
require.NotEmpty(t, evs)
for _, ev := range evs {
require.Equal(t, events.TopicKEVAdded, ev.Topic)
require.Equal(t, kev.Name, ev.Source)
}
}
func TestCollector_EmptyKnownInsertsAll(t *testing.T) {
ftch := &stubFetcher{cat: kev.Catalog{Vulnerabilities: []kev.Vulnerability{
{CveID: "CVE-A", DateAdded: "2026-04-01"},
{CveID: "CVE-B", DateAdded: "2026-04-02"},
{CveID: "CVE-C", DateAdded: "2026-04-03"},
}}}
repo := &stubKEVRepo{known: map[string]bool{}}
emt := &stubKEVEmitter{}
c := kev.NewCollector(kev.CollectorConfig{
Interval: 25 * time.Millisecond,
Fetcher: ftch,
Repo: repo,
Emitter: emt,
State: stubKEVState{},
})
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond)
defer cancel()
_ = c.Run(ctx)
require.Len(t, repo.Inserted(), 3)
require.Len(t, emt.Events(), 3)
}