feat(monitor/collectors/usgs): 1m collector with id-diff and earthquake emit
This commit is contained in:
parent
1d414ff341
commit
2368ed8f18
|
|
@ -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
|
||||
}
|
||||
|
|
@ -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")
|
||||
}
|
||||
Loading…
Reference in New Issue