diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/ransomware/collector.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/ransomware/collector.go new file mode 100644 index 00000000..5343096f --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/ransomware/collector.go @@ -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) +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/ransomware/collector_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/ransomware/collector_test.go new file mode 100644 index 00000000..1651161f --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/ransomware/collector_test.go @@ -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) + } +}