From 2368ed8f18b346c961d80e95559ccb3bc056b17e Mon Sep 17 00:00:00 2001 From: CarterPerez-dev Date: Sat, 2 May 2026 04:22:37 -0400 Subject: [PATCH] feat(monitor/collectors/usgs): 1m collector with id-diff and earthquake emit --- .../internal/collectors/usgs/collector.go | 135 ++++++++++++ .../collectors/usgs/collector_test.go | 194 ++++++++++++++++++ 2 files changed, 329 insertions(+) create mode 100644 PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector.go create mode 100644 PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector_test.go diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector.go new file mode 100644 index 00000000..68cdb5a5 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector.go @@ -0,0 +1,135 @@ +// ©AngelaMos | 2026 +// collector.go + +package usgs + +import ( + "context" + "encoding/json" + "log/slog" + "time" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +const ( + Name = "usgs" + defaultUSGSInterval = time.Minute + defaultLogPlaceLimit = 64 +) + +type Fetcher interface { + Fetch(ctx context.Context) (Feed, error) +} + +type Repository interface { + Upsert(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 = defaultUSGSInterval + } + 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) { + feed, err := c.cfg.Fetcher.Fetch(ctx) + if err != nil { + c.logger.Warn("usgs fetch", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + return + } + + ids := make([]string, 0, len(feed.Features)) + for _, f := range feed.Features { + ids = append(ids, f.ID) + } + known, err := c.cfg.Repo.KnownIDs(ctx, ids) + if err != nil { + c.logger.Warn("usgs known ids", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + return + } + + emitted := int64(0) + now := time.Now().UTC() + for _, f := range feed.Features { + raw, _ := json.Marshal(f) + row := Row{ + ID: f.ID, + OccurredAt: f.Properties.OccurredAt(), + Mag: f.Properties.Mag, + Place: f.Properties.Place, + GeomLat: coord(f.Geometry.Coordinates, 1), + GeomLon: coord(f.Geometry.Coordinates, 0), + DepthKm: coord(f.Geometry.Coordinates, 2), + Payload: raw, + } + if err := c.cfg.Repo.Upsert(ctx, row); err != nil { + c.logger.Warn("usgs upsert", "id", row.ID, "err", err) + continue + } + if known[f.ID] { + continue + } + c.cfg.Emitter.Emit(events.Event{ + Topic: events.TopicEarthquake, + Timestamp: now, + Source: Name, + Payload: json.RawMessage(raw), + }) + emitted++ + } + _ = c.cfg.State.RecordSuccess(ctx, Name, emitted) +} + +func coord(c []float64, i int) float64 { + if i < len(c) { + return c[i] + } + return 0 +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector_test.go new file mode 100644 index 00000000..d1711de0 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/usgs/collector_test.go @@ -0,0 +1,194 @@ +// ©AngelaMos | 2026 +// collector_test.go + +package usgs_test + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/usgs" + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +type fakeFetcher struct { + feed usgs.Feed + err error + mu sync.Mutex + n int +} + +func (f *fakeFetcher) Fetch(_ context.Context) (usgs.Feed, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.n++ + if f.err != nil { + return usgs.Feed{}, f.err + } + return f.feed, nil +} + +type fakeRepo struct { + mu sync.Mutex + known map[string]bool + upserts []usgs.Row +} + +func (r *fakeRepo) Upsert(_ context.Context, row usgs.Row) error { + r.mu.Lock() + defer r.mu.Unlock() + r.upserts = append(r.upserts, row) + if r.known == nil { + r.known = make(map[string]bool) + } + r.known[row.ID] = true + return nil +} + +func (r *fakeRepo) KnownIDs(_ context.Context, ids []string) (map[string]bool, error) { + r.mu.Lock() + defer r.mu.Unlock() + out := make(map[string]bool, len(ids)) + for _, id := range ids { + if r.known[id] { + out[id] = true + } + } + return out, nil +} + +func (r *fakeRepo) Upserts() int { + r.mu.Lock() + defer r.mu.Unlock() + return len(r.upserts) +} + +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) Count() int { + e.mu.Lock() + defer e.mu.Unlock() + return len(e.events) +} + +type recordingState struct { + mu sync.Mutex + successes int + failures int + lastErr string +} + +func (s *recordingState) RecordSuccess(_ context.Context, _ string, _ int64) error { + s.mu.Lock() + defer s.mu.Unlock() + s.successes++ + return nil +} + +func (s *recordingState) RecordError(_ context.Context, _, msg string) error { + s.mu.Lock() + defer s.mu.Unlock() + s.failures++ + s.lastErr = msg + return nil +} + +func TestCollector_TickPersistsAndEmitsNewQuakes(t *testing.T) { + feed := usgs.Feed{ + Type: "FeatureCollection", + Features: []usgs.Feature{ + {ID: "q1", Properties: usgs.Properties{Mag: 4.5, Place: "test 1", Time: time.Now().UnixMilli()}, Geometry: usgs.Geometry{Coordinates: []float64{-120, 49, 5}}}, + {ID: "q2", Properties: usgs.Properties{Mag: 6.5, Place: "test 2", Time: time.Now().UnixMilli()}, Geometry: usgs.Geometry{Coordinates: []float64{140, -30, 10}}}, + {ID: "q3", Properties: usgs.Properties{Mag: 3.0, Place: "test 3", Time: time.Now().UnixMilli()}, Geometry: usgs.Geometry{Coordinates: []float64{0, 0, 1}}}, + }, + } + ftch := &fakeFetcher{feed: feed} + repo := &fakeRepo{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := usgs.NewCollector(usgs.CollectorConfig{ + Interval: 20 * time.Millisecond, + Fetcher: ftch, + Repo: repo, + Emitter: emt, + State: st, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 70*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + require.GreaterOrEqual(t, repo.Upserts(), 3, "should upsert all 3 features at least once") + require.GreaterOrEqual(t, emt.Count(), 3, "should emit 3 new-event events") + for _, ev := range emt.events { + require.Equal(t, events.TopicEarthquake, ev.Topic) + require.Equal(t, usgs.Name, ev.Source) + } + require.Greater(t, st.successes, 0) + require.Equal(t, 0, st.failures) +} + +func TestCollector_FetchErrorRecordsState(t *testing.T) { + ftch := &fakeFetcher{err: errors.New("upstream 503")} + repo := &fakeRepo{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := usgs.NewCollector(usgs.CollectorConfig{ + Interval: 20 * time.Millisecond, + Fetcher: ftch, + Repo: repo, + Emitter: emt, + State: st, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + require.Equal(t, 0, repo.Upserts()) + require.Equal(t, 0, emt.Count()) + require.Greater(t, st.failures, 0) + require.Contains(t, st.lastErr, "upstream 503") +} + +func TestCollector_KnownQuakesNotReEmitted(t *testing.T) { + feed := usgs.Feed{ + Features: []usgs.Feature{ + {ID: "qx", Properties: usgs.Properties{Mag: 4.5, Time: time.Now().UnixMilli()}, Geometry: usgs.Geometry{Coordinates: []float64{0, 0, 1}}}, + }, + } + ftch := &fakeFetcher{feed: feed} + repo := &fakeRepo{known: map[string]bool{"qx": true}} + emt := &fakeEmitter{} + st := &recordingState{} + + c := usgs.NewCollector(usgs.CollectorConfig{ + Interval: 20 * time.Millisecond, + Fetcher: ftch, + Repo: repo, + Emitter: emt, + State: st, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + require.Equal(t, 0, emt.Count(), "known quake should not re-emit") +}