feat(monitor/collectors/cfradar): 5m collector with diff-on-ID emitting outages + hijacks
This commit is contained in:
parent
428610b663
commit
ee358e38f1
|
|
@ -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
|
||||
}
|
||||
|
|
@ -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")
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue