diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/client.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/client.go new file mode 100644 index 00000000..5ba0c445 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/client.go @@ -0,0 +1,75 @@ +// ©AngelaMos | 2026 +// client.go + +package wikipedia + +import ( + "context" + "fmt" + "net/url" + "time" + + "golang.org/x/time/rate" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/httpx" +) + +const ( + defaultWikiBaseURL = "https://en.wikipedia.org" + pathAPI = "/w/api.php" + defaultWikiRate = 200 * time.Millisecond + defaultWikiBurst = 5 + defaultWikiBudget = 5 + defaultWikiBreaker = 60 * time.Second +) + +type ClientConfig struct { + BaseURL string +} + +type Client struct { + hx *httpx.Client +} + +func NewClient(cfg ClientConfig) *Client { + if cfg.BaseURL == "" { + cfg.BaseURL = defaultWikiBaseURL + } + return &Client{ + hx: httpx.New(httpx.Config{ + Name: "wikipedia", + BaseURL: cfg.BaseURL, + Rate: rate.Every(defaultWikiRate), + Burst: defaultWikiBurst, + ConsecutiveFailureBudget: defaultWikiBudget, + BreakerTimeout: defaultWikiBreaker, + }), + } +} + +func (c *Client) Fetch(ctx context.Context) (Response, error) { + q := url.Values{} + q.Set("action", "parse") + q.Set("page", "Template:In_the_news") + q.Set("prop", "text|revid") + q.Set("format", "json") + + resp, err := c.hx.Get(ctx, pathAPI, q) + if err != nil { + return Response{}, fmt.Errorf("fetch wikipedia ITN: %w", err) + } + defer func() { _ = resp.Body.Close() }() + + body := make([]byte, 0, 32*1024) + buf := make([]byte, 4*1024) + for { + n, rerr := resp.Body.Read(buf) + if n > 0 { + body = append(body, buf[:n]...) + } + if rerr != nil { + break + } + } + return DecodeResponse(body) +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/client_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/client_test.go new file mode 100644 index 00000000..6732c28f --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/client_test.go @@ -0,0 +1,40 @@ +// ©AngelaMos | 2026 +// client_test.go + +package wikipedia_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/wikipedia" +) + +func TestClient_FetchDecodesITNResponse(t *testing.T) { + body, err := os.ReadFile("testdata/itn_response.json") + require.NoError(t, err) + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + require.Equal(t, "/w/api.php", r.URL.Path) + require.Equal(t, "parse", r.URL.Query().Get("action")) + require.Equal(t, "Template:In_the_news", r.URL.Query().Get("page")) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(body) + })) + defer srv.Close() + + c := wikipedia.NewClient(wikipedia.ClientConfig{BaseURL: srv.URL}) + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + resp, err := c.Fetch(ctx) + require.NoError(t, err) + require.NotZero(t, resp.RevID) + require.NotEmpty(t, resp.HTML) +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/collector.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/collector.go new file mode 100644 index 00000000..2262ed88 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/collector.go @@ -0,0 +1,138 @@ +// ©AngelaMos | 2026 +// collector.go + +package wikipedia + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "log/slog" + "time" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +const ( + Name = "wikipedia" + defaultWikiInterval = 5 * time.Minute + idHashBytes = 16 +) + +type Fetcher interface { + Fetch(ctx context.Context) (Response, error) +} + +type Repository interface { + RememberRevID(ctx context.Context, revID int64) error + LastRevID(ctx context.Context) (int64, bool, error) + Insert(ctx context.Context, e Entry) 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 = defaultWikiInterval + } + 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) { + resp, err := c.cfg.Fetcher.Fetch(ctx) + if err != nil { + c.logger.Warn("wikipedia fetch", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + return + } + + last, found, err := c.cfg.Repo.LastRevID(ctx) + if err != nil { + c.logger.Warn("wikipedia revid lookup", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + return + } + if found && last == resp.RevID { + _ = c.cfg.State.RecordSuccess(ctx, Name, 0) + return + } + + entries := ParseEntries(resp.HTML) + now := time.Now().UTC() + emitted := int64(0) + for _, e := range entries { + id := entryID(e) + body, _ := json.Marshal(map[string]any{ + "text": e.Text, + "slug": e.ArticleSlug, + }) + entry := Entry{ + ID: id, + Headline: e.Text, + OccurredAt: now, + Payload: body, + } + if err := c.cfg.Repo.Insert(ctx, entry); err != nil { + c.logger.Warn("wikipedia insert", "id", id, "err", err) + continue + } + c.cfg.Emitter.Emit(events.Event{ + Topic: events.TopicWikipediaITN, + Timestamp: now, + Source: Name, + Payload: json.RawMessage(body), + }) + emitted++ + } + + if err := c.cfg.Repo.RememberRevID(ctx, resp.RevID); err != nil { + c.logger.Warn("wikipedia remember revid", "err", err) + } + _ = c.cfg.State.RecordSuccess(ctx, Name, emitted) +} + +func entryID(e ITNEntry) string { + h := sha256.Sum256([]byte(e.Text + "|" + e.ArticleSlug)) + return hex.EncodeToString(h[:idHashBytes]) +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/collector_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/collector_test.go new file mode 100644 index 00000000..628fd9a8 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/collector_test.go @@ -0,0 +1,178 @@ +// ©AngelaMos | 2026 +// collector_test.go + +package wikipedia_test + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/wikipedia" + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +type fakeFetcher struct { + resp wikipedia.Response + err error +} + +func (f *fakeFetcher) Fetch(_ context.Context) (wikipedia.Response, error) { + return f.resp, f.err +} + +type fakeRepo struct { + mu sync.Mutex + revID int64 + revKnown bool + inserts []wikipedia.Entry + saveCalls int +} + +func (r *fakeRepo) RememberRevID(_ context.Context, revID int64) error { + r.mu.Lock() + defer r.mu.Unlock() + r.revID = revID + r.revKnown = true + r.saveCalls++ + return nil +} + +func (r *fakeRepo) LastRevID(_ context.Context) (int64, bool, error) { + r.mu.Lock() + defer r.mu.Unlock() + return r.revID, r.revKnown, nil +} + +func (r *fakeRepo) Insert(_ context.Context, e wikipedia.Entry) error { + r.mu.Lock() + defer r.mu.Unlock() + r.inserts = append(r.inserts, e) + 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_NewRevidInsertsAndEmits(t *testing.T) { + resp := wikipedia.Response{ + RevID: 999, + HTML: ``, + } + ftch := &fakeFetcher{resp: resp} + repo := &fakeRepo{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := wikipedia.NewCollector(wikipedia.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.GreaterOrEqual(t, repo.Inserts(), 2) + require.GreaterOrEqual(t, emt.Count(), 2) + for _, ev := range emt.events { + require.Equal(t, events.TopicWikipediaITN, ev.Topic) + } + require.Greater(t, st.successes, 0) +} + +func TestCollector_RevIDUnchangedSkipsInsert(t *testing.T) { + resp := wikipedia.Response{ + RevID: 555, + HTML: ``, + } + ftch := &fakeFetcher{resp: resp} + repo := &fakeRepo{revID: 555, revKnown: true} + emt := &fakeEmitter{} + st := &recordingState{} + + c := wikipedia.NewCollector(wikipedia.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.Inserts()) + require.Equal(t, 0, emt.Count()) + require.Greater(t, st.successes, 0) +} + +func TestCollector_FetchErrorRecordsState(t *testing.T) { + ftch := &fakeFetcher{err: errors.New("upstream 503")} + repo := &fakeRepo{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := wikipedia.NewCollector(wikipedia.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.Inserts()) + require.Greater(t, st.failures, 0) +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/repo.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/repo.go new file mode 100644 index 00000000..6f54aba4 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/wikipedia/repo.go @@ -0,0 +1,72 @@ +// ©AngelaMos | 2026 +// repo.go + +package wikipedia + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "strconv" + "time" + + "github.com/jmoiron/sqlx" + "github.com/redis/go-redis/v9" +) + +const ( + keyRevID = "state:wiki_itn:revid" + sourceWikiITN = "wiki_itn" +) + +type Entry struct { + ID string + Headline string + OccurredAt time.Time + Payload json.RawMessage +} + +type Repo struct { + db *sqlx.DB + rdb *redis.Client +} + +func NewRepo(db *sqlx.DB, rdb *redis.Client) *Repo { + return &Repo{db: db, rdb: rdb} +} + +func (r *Repo) RememberRevID(ctx context.Context, revID int64) error { + if err := r.rdb.Set(ctx, keyRevID, strconv.FormatInt(revID, 10), 0).Err(); err != nil { + return fmt.Errorf("save wiki revid: %w", err) + } + return nil +} + +func (r *Repo) LastRevID(ctx context.Context) (int64, bool, error) { + s, err := r.rdb.Get(ctx, keyRevID).Result() + if errors.Is(err, redis.Nil) { + return 0, false, nil + } + if err != nil { + return 0, false, fmt.Errorf("load wiki revid: %w", err) + } + rev, err := strconv.ParseInt(s, 10, 64) + if err != nil { + return 0, false, fmt.Errorf("parse wiki revid %q: %w", s, err) + } + return rev, true, nil +} + +func (r *Repo) Insert(ctx context.Context, e Entry) 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`, + e.ID, sourceWikiITN, e.OccurredAt, e.Headline, []byte(e.Payload), + ) + if err != nil { + return fmt.Errorf("insert wiki event %s: %w", e.ID, err) + } + return nil +}