feat(monitor/collectors/gdelt): 15m collector with per-theme rolling z-score spike emit
This commit is contained in:
parent
066e2debc2
commit
317fadeeef
|
|
@ -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
|
||||
}
|
||||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
|
|
@ -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])
|
||||
}
|
||||
|
|
@ -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)
|
||||
}
|
||||
|
|
@ -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
|
||||
}
|
||||
|
|
@ -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}
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
Loading…
Reference in New Issue