diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/client.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/client.go new file mode 100644 index 00000000..ec2bd0ff --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/client.go @@ -0,0 +1,123 @@ +// ©AngelaMos | 2026 +// client.go + +package gdelt + +import ( + "context" + "encoding/json" + "fmt" + "net/url" + "strings" + "time" + + "golang.org/x/time/rate" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/httpx" +) + +const ( + defaultGDELTBaseURL = "https://api.gdeltproject.org" + pathDoc = "/api/v2/doc/doc" + defaultGDELTRate = 500 * time.Millisecond + defaultGDELTBurst = 2 + defaultGDELTBudget = 5 + defaultGDELTBreaker = 60 * time.Second + + gdeltDateFormat = "20060102T150405Z" +) + +var DefaultThemes = []string{ + "NATURAL_DISASTER", + "ARMEDCONFLICT", + "DISEASE_OUTBREAK", + "ECON_BANKRUPTCY", + "TERROR", +} + +type ClientConfig struct { + BaseURL string +} + +type Client struct { + hx *httpx.Client +} + +func NewClient(cfg ClientConfig) *Client { + if cfg.BaseURL == "" { + cfg.BaseURL = defaultGDELTBaseURL + } + return &Client{ + hx: httpx.New(httpx.Config{ + Name: "gdelt", + BaseURL: cfg.BaseURL, + Rate: rate.Every(defaultGDELTRate), + Burst: defaultGDELTBurst, + ConsecutiveFailureBudget: defaultGDELTBudget, + BreakerTimeout: defaultGDELTBreaker, + }), + } +} + +type ThemeBucket struct { + Theme string + Time time.Time + Count int +} + +type rawTimeline struct { + Timeline []struct { + Data []struct { + Date string `json:"date"` + Value int `json:"value"` + } `json:"data"` + } `json:"timeline"` +} + +func (c *Client) FetchTheme(ctx context.Context, theme string) ([]ThemeBucket, error) { + q := url.Values{} + q.Set("query", "theme:"+theme) + q.Set("mode", "timelinevolinfo") + q.Set("TIMELINESMOOTH", "5") + q.Set("format", "json") + q.Set("timespan", "1d") + + resp, err := c.hx.Get(ctx, pathDoc, q) + if err != nil { + return nil, fmt.Errorf("fetch gdelt theme %s: %w", theme, err) + } + defer func() { _ = resp.Body.Close() }() + + body := strings.Builder{} + buf := make([]byte, 4096) + for { + n, rerr := resp.Body.Read(buf) + if n > 0 { + body.Write(buf[:n]) + } + if rerr != nil { + break + } + } + + var raw rawTimeline + if err := json.Unmarshal([]byte(body.String()), &raw); err != nil { + return nil, fmt.Errorf("decode gdelt theme %s: %w", theme, err) + } + if len(raw.Timeline) == 0 { + return nil, nil + } + out := make([]ThemeBucket, 0, len(raw.Timeline[0].Data)) + for _, d := range raw.Timeline[0].Data { + ts, err := time.Parse(gdeltDateFormat, d.Date) + if err != nil { + continue + } + out = append(out, ThemeBucket{ + Theme: theme, + Time: ts.UTC(), + Count: d.Value, + }) + } + return out, nil +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/client_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/client_test.go new file mode 100644 index 00000000..0f7e4c26 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/client_test.go @@ -0,0 +1,43 @@ +// ©AngelaMos | 2026 +// client_test.go + +package gdelt_test + +import ( + "context" + "net/http" + "net/http/httptest" + "os" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/gdelt" +) + +func TestClient_FetchThemeDecodesBuckets(t *testing.T) { + body, err := os.ReadFile("testdata/timelinevolinfo.json") + require.NoError(t, err) + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + require.Equal(t, "/api/v2/doc/doc", r.URL.Path) + require.Equal(t, "theme:NATURAL_DISASTER", r.URL.Query().Get("query")) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(body) + })) + defer srv.Close() + + c := gdelt.NewClient(gdelt.ClientConfig{BaseURL: srv.URL}) + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + buckets, err := c.FetchTheme(ctx, "NATURAL_DISASTER") + require.NoError(t, err) + require.GreaterOrEqual(t, len(buckets), 1) + for _, b := range buckets { + require.False(t, b.Time.IsZero()) + require.Greater(t, b.Count, 0) + require.Equal(t, "NATURAL_DISASTER", b.Theme) + } +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/collector.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/collector.go new file mode 100644 index 00000000..cf1d491b --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/collector.go @@ -0,0 +1,163 @@ +// ©AngelaMos | 2026 +// collector.go + +package gdelt + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "log/slog" + "time" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +const ( + Name = "gdelt" + defaultGDELTInterval = 15 * time.Minute + defaultBaselineCap = 96 + zScoreSpikeThreshold = 3.0 + idHashBytes = 16 +) + +type Fetcher interface { + FetchTheme(ctx context.Context, theme string) ([]ThemeBucket, error) +} + +type Repository interface { + Insert(ctx context.Context, row SpikeRow) 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 + Themes []string + BaselineCap int + Fetcher Fetcher + Repo Repository + Emitter Emitter + State StateRecorder + Logger *slog.Logger +} + +type Collector struct { + cfg CollectorConfig + logger *slog.Logger + baselines map[string]*ThemeState + emitted map[string]bool +} + +func NewCollector(cfg CollectorConfig) *Collector { + if cfg.Interval <= 0 { + cfg.Interval = defaultGDELTInterval + } + if cfg.BaselineCap <= 0 { + cfg.BaselineCap = defaultBaselineCap + } + if len(cfg.Themes) == 0 { + cfg.Themes = DefaultThemes + } + if cfg.Logger == nil { + cfg.Logger = slog.Default() + } + c := &Collector{ + cfg: cfg, + logger: cfg.Logger, + baselines: make(map[string]*ThemeState, len(cfg.Themes)), + emitted: make(map[string]bool), + } + for _, t := range cfg.Themes { + c.baselines[t] = NewThemeState(cfg.BaselineCap) + } + return c +} + +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) { + hadError := false + emitted := int64(0) + + for _, theme := range c.cfg.Themes { + buckets, err := c.cfg.Fetcher.FetchTheme(ctx, theme) + if err != nil { + c.logger.Warn("gdelt fetch", "theme", theme, "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + hadError = true + continue + } + baseline := c.baselines[theme] + for _, b := range buckets { + z := baseline.ZScore(b.Count) + baseline.Push(Bucket{Score: b.Time.UnixMilli(), Count: b.Count}) + + if z <= zScoreSpikeThreshold { + continue + } + id := spikeID(theme, b.Time) + if c.emitted[id] { + continue + } + c.emitted[id] = true + + payload, _ := json.Marshal(map[string]any{ + "theme": theme, + "time": b.Time, + "count": b.Count, + "zscore": z, + }) + row := SpikeRow{ + ID: id, + Theme: theme, + OccurredAt: b.Time, + Headline: fmt.Sprintf("Theme spike: %s (z=%.2f, count=%d)", theme, z, b.Count), + Payload: payload, + } + if err := c.cfg.Repo.Insert(ctx, row); err != nil { + c.logger.Warn("gdelt insert", "id", id, "err", err) + continue + } + c.cfg.Emitter.Emit(events.Event{ + Topic: events.TopicGDELTSpike, + Timestamp: b.Time, + Source: Name, + Payload: json.RawMessage(payload), + }) + emitted++ + } + } + + if !hadError { + _ = c.cfg.State.RecordSuccess(ctx, Name, emitted) + } +} + +func spikeID(theme string, t time.Time) string { + h := sha256.Sum256([]byte(theme + "|" + t.UTC().Format(time.RFC3339))) + return hex.EncodeToString(h[:idHashBytes]) +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/collector_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/collector_test.go new file mode 100644 index 00000000..a599adbf --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/collector_test.go @@ -0,0 +1,181 @@ +// ©AngelaMos | 2026 +// collector_test.go + +package gdelt_test + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/gdelt" + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +type fakeFetcher struct { + mu sync.Mutex + calls int + buckets map[string][]gdelt.ThemeBucket + err error +} + +func (f *fakeFetcher) FetchTheme(_ context.Context, theme string) ([]gdelt.ThemeBucket, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.calls++ + if f.err != nil { + return nil, f.err + } + return f.buckets[theme], nil +} + +type fakeRepo struct { + mu sync.Mutex + inserts []gdelt.SpikeRow +} + +func (r *fakeRepo) Insert(_ context.Context, row gdelt.SpikeRow) error { + r.mu.Lock() + defer r.mu.Unlock() + r.inserts = append(r.inserts, row) + return nil +} + +func (r *fakeRepo) Inserts() int { + r.mu.Lock() + defer r.mu.Unlock() + return len(r.inserts) +} + +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 +} + +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, _, _ string) error { + s.mu.Lock() + defer s.mu.Unlock() + s.failures++ + return nil +} + +func TestCollector_NoBaselineNoSpike(t *testing.T) { + now := time.Now().UTC() + buckets := []gdelt.ThemeBucket{ + {Theme: "X", Time: now, Count: 5000}, + } + ftch := &fakeFetcher{buckets: map[string][]gdelt.ThemeBucket{"X": buckets}} + repo := &fakeRepo{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := gdelt.NewCollector(gdelt.CollectorConfig{ + Interval: 20 * time.Millisecond, + Themes: []string{"X"}, + BaselineCap: 8, + 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.Inserts(), "first observation has no baseline → cannot detect spike") + require.Equal(t, 0, emt.Count()) + require.Greater(t, st.successes, 0) +} + +func TestCollector_StableBaselinePlusSpikeEmitsOnce(t *testing.T) { + base := time.Now().UTC().Truncate(15 * time.Minute) + stable := make([]gdelt.ThemeBucket, 0, 10) + for i := 0; i < 10; i++ { + stable = append(stable, gdelt.ThemeBucket{ + Theme: "X", + Time: base.Add(time.Duration(i) * 15 * time.Minute), + Count: 100 + i, + }) + } + spikeBucket := gdelt.ThemeBucket{Theme: "X", Time: base.Add(15 * 15 * time.Minute), Count: 5000} + buckets := append(stable, spikeBucket) + + ftch := &fakeFetcher{buckets: map[string][]gdelt.ThemeBucket{"X": buckets}} + repo := &fakeRepo{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := gdelt.NewCollector(gdelt.CollectorConfig{ + Interval: 20 * time.Millisecond, + Themes: []string{"X"}, + BaselineCap: 8, + Fetcher: ftch, + Repo: repo, + Emitter: emt, + State: st, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 80*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + require.GreaterOrEqual(t, repo.Inserts(), 1, "spike must be inserted") + require.GreaterOrEqual(t, emt.Count(), 1) + for _, ev := range emt.events { + require.Equal(t, events.TopicGDELTSpike, ev.Topic) + } +} + +func TestCollector_FetchErrorsRecordsState(t *testing.T) { + ftch := &fakeFetcher{err: errors.New("upstream 503")} + repo := &fakeRepo{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := gdelt.NewCollector(gdelt.CollectorConfig{ + Interval: 20 * time.Millisecond, + Themes: []string{"X"}, + BaselineCap: 8, + Fetcher: ftch, + Repo: repo, + Emitter: emt, + State: st, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + require.Equal(t, 0, repo.Inserts()) + require.Greater(t, st.failures, 0) +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/repo.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/repo.go new file mode 100644 index 00000000..93a09c51 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/repo.go @@ -0,0 +1,44 @@ +// ©AngelaMos | 2026 +// repo.go + +package gdelt + +import ( + "context" + "encoding/json" + "fmt" + "time" + + "github.com/jmoiron/sqlx" +) + +const ( + sourceGDELTSpike = "gdelt_spike" +) + +type SpikeRow struct { + ID string + Theme string + OccurredAt time.Time + Headline string + Payload json.RawMessage +} + +type Repo struct { + db *sqlx.DB +} + +func NewRepo(db *sqlx.DB) *Repo { return &Repo{db: db} } + +func (r *Repo) Insert(ctx context.Context, row SpikeRow) error { + _, err := r.db.ExecContext(ctx, ` + INSERT INTO world_events (id, source, occurred_at, headline, payload) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (id) DO NOTHING`, + row.ID, sourceGDELTSpike, row.OccurredAt, row.Headline, []byte(row.Payload), + ) + if err != nil { + return fmt.Errorf("insert gdelt spike %s: %w", row.ID, err) + } + return nil +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/testdata/timelinevolinfo.json b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/testdata/timelinevolinfo.json new file mode 100644 index 00000000..6aef353b --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/gdelt/testdata/timelinevolinfo.json @@ -0,0 +1,16 @@ +{ + "timeline": [ + { + "data": [ + {"date": "20260501T120000Z", "value": 100}, + {"date": "20260501T121500Z", "value": 105}, + {"date": "20260501T123000Z", "value": 102}, + {"date": "20260501T124500Z", "value": 108}, + {"date": "20260501T130000Z", "value": 112}, + {"date": "20260501T131500Z", "value": 98}, + {"date": "20260501T133000Z", "value": 103}, + {"date": "20260501T134500Z", "value": 107} + ] + } + ] +}