diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/cfradar/collector.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/cfradar/collector.go new file mode 100644 index 00000000..d8c9c09a --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/cfradar/collector.go @@ -0,0 +1,212 @@ +// ©AngelaMos | 2026 +// collector.go + +package cfradar + +import ( + "context" + "encoding/json" + "log/slog" + "time" + + "github.com/lib/pq" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +const ( + Name = "cfradar" + defaultCollectorInterval = 5 * time.Minute + defaultCollectorConfidence = 7 +) + +type Fetcher interface { + FetchOutages(ctx context.Context) (OutageResultBody, error) + FetchHijacks(ctx context.Context, minConfidence int) (HijackBody, error) +} + +type Repository interface { + UpsertOutage(ctx context.Context, o OutageRow) error + UpsertHijack(ctx context.Context, h HijackRow) error + KnownOutageIDs(ctx context.Context, ids []string) (map[string]bool, error) + KnownHijackIDs(ctx context.Context, ids []int64) (map[int64]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 + MinConfidence int + 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 = defaultCollectorInterval + } + if cfg.MinConfidence <= 0 { + cfg.MinConfidence = defaultCollectorConfidence + } + 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) { + emitted := int64(0) + hadError := false + + if n, err := c.tickOutages(ctx); err != nil { + c.logger.Warn("cfradar outages tick failed", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + hadError = true + } else { + emitted += n + } + + if n, err := c.tickHijacks(ctx); err != nil { + c.logger.Warn("cfradar hijacks tick failed", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + hadError = true + } else { + emitted += n + } + + if !hadError { + _ = c.cfg.State.RecordSuccess(ctx, Name, emitted) + } +} + +func (c *Collector) tickOutages(ctx context.Context) (int64, error) { + body, err := c.cfg.Fetcher.FetchOutages(ctx) + if err != nil { + return 0, err + } + + ids := make([]string, 0, len(body.Annotations)) + for _, a := range body.Annotations { + ids = append(ids, a.ID) + } + + known, err := c.cfg.Repo.KnownOutageIDs(ctx, ids) + if err != nil { + return 0, err + } + + now := time.Now().UTC() + emitted := int64(0) + for _, a := range body.Annotations { + if known[a.ID] { + continue + } + rawBytes, _ := json.Marshal(a) + raw := json.RawMessage(rawBytes) + row := OutageRow{ + ID: a.ID, + StartedAt: a.StartDate, + EndedAt: a.EndDate, + Locations: pq.StringArray(a.Locations), + ASNs: pq.Int32Array(a.ASNs), + Cause: a.Reason, + OutageType: a.OutageType, + Payload: raw, + } + if uerr := c.cfg.Repo.UpsertOutage(ctx, row); uerr != nil { + c.logger.Warn("upsert outage", "id", a.ID, "err", uerr) + continue + } + c.cfg.Emitter.Emit(events.Event{ + Topic: events.TopicInternetOutage, + Timestamp: now, + Source: Name, + Payload: raw, + }) + emitted++ + } + return emitted, nil +} + +func (c *Collector) tickHijacks(ctx context.Context) (int64, error) { + body, err := c.cfg.Fetcher.FetchHijacks(ctx, c.cfg.MinConfidence) + if err != nil { + return 0, err + } + + ids := make([]int64, 0, len(body.Events)) + for _, e := range body.Events { + ids = append(ids, e.ID) + } + + known, err := c.cfg.Repo.KnownHijackIDs(ctx, ids) + if err != nil { + return 0, err + } + + now := time.Now().UTC() + emitted := int64(0) + for _, e := range body.Events { + if known[e.ID] { + continue + } + rawBytes, _ := json.Marshal(e) + raw := json.RawMessage(rawBytes) + row := HijackRow{ + ID: e.ID, + DetectedAt: e.DetectedAt, + StartedAt: e.StartedAt, + DurationSec: e.DurationSec, + Confidence: e.Confidence, + HijackerASN: e.HijackerASN, + VictimASNs: pq.Int32Array(e.VictimASNs), + Prefixes: e.Prefixes, + Payload: raw, + } + if uerr := c.cfg.Repo.UpsertHijack(ctx, row); uerr != nil { + c.logger.Warn("upsert hijack", "id", e.ID, "err", uerr) + continue + } + c.cfg.Emitter.Emit(events.Event{ + Topic: events.TopicBGPHijack, + Timestamp: now, + Source: Name, + Payload: raw, + }) + emitted++ + } + return emitted, nil +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/cfradar/collector_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/cfradar/collector_test.go new file mode 100644 index 00000000..c0f93853 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/cfradar/collector_test.go @@ -0,0 +1,191 @@ +// ©AngelaMos | 2026 +// collector_test.go + +package cfradar_test + +import ( + "context" + "encoding/json" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/cfradar" + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +type fakeFetcher struct { + outages cfradar.OutageResultBody + hijacks cfradar.HijackBody +} + +func (f *fakeFetcher) FetchOutages(context.Context) (cfradar.OutageResultBody, error) { + return f.outages, nil +} + +func (f *fakeFetcher) FetchHijacks(context.Context, int) (cfradar.HijackBody, error) { + return f.hijacks, nil +} + +type fakeRepo struct { + mu sync.Mutex + outageUpserts int + hijackUpserts int + knownOutages map[string]bool + knownHijacks map[int64]bool +} + +func (r *fakeRepo) UpsertOutage(_ context.Context, _ cfradar.OutageRow) error { + r.mu.Lock() + defer r.mu.Unlock() + r.outageUpserts++ + return nil +} + +func (r *fakeRepo) UpsertHijack(_ context.Context, _ cfradar.HijackRow) error { + r.mu.Lock() + defer r.mu.Unlock() + r.hijackUpserts++ + return nil +} + +func (r *fakeRepo) KnownOutageIDs(_ context.Context, ids []string) (map[string]bool, error) { + out := make(map[string]bool) + for _, id := range ids { + if r.knownOutages[id] { + out[id] = true + } + } + return out, nil +} + +func (r *fakeRepo) KnownHijackIDs(_ context.Context, ids []int64) (map[int64]bool, error) { + out := make(map[int64]bool) + for _, id := range ids { + if r.knownHijacks[id] { + out[id] = true + } + } + return out, nil +} + +type fakeEmitter struct { + mu sync.Mutex + events []events.Event +} + +func (e *fakeEmitter) Emit(ev events.Event) { + e.mu.Lock() + defer e.mu.Unlock() + e.events = append(e.events, ev) +} + +func (e *fakeEmitter) Events() []events.Event { + e.mu.Lock() + defer e.mu.Unlock() + out := make([]events.Event, len(e.events)) + copy(out, e.events) + return out +} + +type noopState struct{} + +func (noopState) RecordSuccess(context.Context, string, int64) error { return nil } +func (noopState) RecordError(context.Context, string, string) error { return nil } + +func TestCollector_OnlyEmitsNetNew(t *testing.T) { + now := time.Now().UTC() + ftch := &fakeFetcher{ + outages: cfradar.OutageResultBody{Annotations: []cfradar.OutageAnnotation{ + {ID: "out-known", StartDate: now}, + {ID: "out-new", StartDate: now}, + }}, + hijacks: cfradar.HijackBody{Events: []cfradar.HijackEvent{ + {ID: 100, DetectedAt: now, StartedAt: now, Confidence: 9}, + {ID: 200, DetectedAt: now, StartedAt: now, Confidence: 8}, + }}, + } + repo := &fakeRepo{ + knownOutages: map[string]bool{"out-known": true}, + knownHijacks: map[int64]bool{200: true}, + } + emt := &fakeEmitter{} + + c := cfradar.NewCollector(cfradar.CollectorConfig{ + Interval: 30 * time.Millisecond, + MinConfidence: 7, + Fetcher: ftch, + Repo: repo, + Emitter: emt, + State: noopState{}, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + evs := emt.Events() + + var outageEvents, hijackEvents int + for _, ev := range evs { + switch ev.Topic { + case events.TopicInternetOutage: + outageEvents++ + body, _ := json.Marshal(ev.Payload) + require.Contains(t, string(body), "out-new") + require.NotContains(t, string(body), "out-known") + case events.TopicBGPHijack: + hijackEvents++ + body, _ := json.Marshal(ev.Payload) + require.Contains(t, string(body), "100") + } + } + require.GreaterOrEqual(t, outageEvents, 1) + require.GreaterOrEqual(t, hijackEvents, 1) +} + +func TestCollector_RepeatedTickIsIdempotent(t *testing.T) { + now := time.Now().UTC() + ftch := &fakeFetcher{ + outages: cfradar.OutageResultBody{Annotations: []cfradar.OutageAnnotation{ + {ID: "out-x", StartDate: now}, + }}, + } + emit := &fakeEmitter{} + + known := map[string]bool{} + repo := &fakeRepo{knownOutages: known} + + c := cfradar.NewCollector(cfradar.CollectorConfig{ + Interval: 20 * time.Millisecond, + MinConfidence: 7, + Fetcher: ftch, + Repo: repo, + Emitter: emit, + State: noopState{}, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 25*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + known["out-x"] = true + + emit2 := &fakeEmitter{} + c2 := cfradar.NewCollector(cfradar.CollectorConfig{ + Interval: 20 * time.Millisecond, + MinConfidence: 7, + Fetcher: ftch, + Repo: repo, + Emitter: emit2, + State: noopState{}, + }) + ctx2, cancel2 := context.WithTimeout(context.Background(), 25*time.Millisecond) + defer cancel2() + _ = c2.Run(ctx2) + + for _, ev := range emit2.Events() { + require.NotEqual(t, events.TopicInternetOutage, ev.Topic, "should not re-emit known outage") + } +}