feat(monitor/collectors/ransomware): 15m collector with hash-diff and victim emit

This commit is contained in:
CarterPerez-dev 2026-05-01 22:48:09 -04:00
parent 3d71fd8c84
commit 894ecc8fe5
2 changed files with 251 additions and 0 deletions

View File

@ -0,0 +1,133 @@
// ©AngelaMos | 2026
// collector.go
package ransomware
import (
"context"
"encoding/json"
"log/slog"
"time"
"github.com/carterperez-dev/monitor-the-situation/backend/internal/events"
)
const (
Name = "ransomware"
defaultRansomCadence = 15 * time.Minute
)
type Fetcher interface {
FetchRecent(ctx context.Context) ([]Victim, 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 = defaultRansomCadence
}
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) {
vs, err := c.cfg.Fetcher.FetchRecent(ctx)
if err != nil {
c.logger.Warn("ransomware fetch", "err", err)
_ = c.cfg.State.RecordError(ctx, Name, err.Error())
return
}
ids := make([]string, 0, len(vs))
for _, v := range vs {
ids = append(ids, v.ID())
}
known, err := c.cfg.Repo.KnownIDs(ctx, ids)
if err != nil {
c.logger.Warn("ransomware known ids", "err", err)
_ = c.cfg.State.RecordError(ctx, Name, err.Error())
return
}
now := time.Now().UTC()
emitted := int64(0)
for _, v := range vs {
id := v.ID()
if known[id] {
continue
}
raw, _ := json.Marshal(v)
row := Row{
ID: id,
PostTitle: v.PostTitle,
GroupName: v.GroupName,
DiscoveredAt: v.Discovered,
Country: v.Country,
Sector: v.Activity,
Payload: raw,
}
if err := c.cfg.Repo.Insert(ctx, row); err != nil {
c.logger.Warn("ransomware insert", "id", id, "err", err)
continue
}
c.cfg.Emitter.Emit(events.Event{
Topic: events.TopicRansomwareVictim,
Timestamp: now,
Source: Name,
Payload: json.RawMessage(raw),
})
emitted++
}
_ = c.cfg.State.RecordSuccess(ctx, Name, emitted)
}

View File

@ -0,0 +1,118 @@
// ©AngelaMos | 2026
// collector_test.go
package ransomware_test
import (
"context"
"sync"
"testing"
"time"
"github.com/stretchr/testify/require"
"github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/ransomware"
"github.com/carterperez-dev/monitor-the-situation/backend/internal/events"
)
type stubFetcher struct {
victims []ransomware.Victim
}
func (s *stubFetcher) FetchRecent(context.Context) ([]ransomware.Victim, error) {
return s.victims, nil
}
type stubRansomRepo struct {
mu sync.Mutex
inserted []string
known map[string]bool
}
func (r *stubRansomRepo) Insert(_ context.Context, row ransomware.Row) error {
r.mu.Lock()
defer r.mu.Unlock()
r.inserted = append(r.inserted, row.ID)
if r.known == nil {
r.known = map[string]bool{}
}
r.known[row.ID] = true
return nil
}
func (r *stubRansomRepo) 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 *stubRansomRepo) Inserted() []string {
r.mu.Lock()
defer r.mu.Unlock()
out := make([]string, len(r.inserted))
copy(out, r.inserted)
return out
}
type stubRansomEmitter struct {
mu sync.Mutex
events []events.Event
}
func (e *stubRansomEmitter) Emit(ev events.Event) {
e.mu.Lock()
defer e.mu.Unlock()
e.events = append(e.events, ev)
}
func (e *stubRansomEmitter) Events() []events.Event {
e.mu.Lock()
defer e.mu.Unlock()
out := make([]events.Event, len(e.events))
copy(out, e.events)
return out
}
type stubRansomState struct{}
func (stubRansomState) RecordSuccess(context.Context, string, int64) error { return nil }
func (stubRansomState) RecordError(context.Context, string, string) error { return nil }
func TestCollector_OnlyEmitsNewVictims(t *testing.T) {
now := time.Now().UTC()
known := ransomware.Victim{PostTitle: "Old", GroupName: "lockbit", Discovered: now.Add(-time.Hour)}
new1 := ransomware.Victim{PostTitle: "Acme", GroupName: "blackcat", Discovered: now}
new2 := ransomware.Victim{PostTitle: "Banco", GroupName: "play", Discovered: now}
ftch := &stubFetcher{victims: []ransomware.Victim{known, new1, new2}}
repo := &stubRansomRepo{known: map[string]bool{known.ID(): true}}
emt := &stubRansomEmitter{}
c := ransomware.NewCollector(ransomware.CollectorConfig{
Interval: 30 * time.Millisecond,
Fetcher: ftch,
Repo: repo,
Emitter: emt,
State: stubRansomState{},
})
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
defer cancel()
_ = c.Run(ctx)
inserted := repo.Inserted()
require.Len(t, inserted, 2)
evs := emt.Events()
require.Len(t, evs, 2)
for _, ev := range evs {
require.Equal(t, events.TopicRansomwareVictim, ev.Topic)
require.Equal(t, ransomware.Name, ev.Source)
}
}