diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/swpc/collector.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/swpc/collector.go new file mode 100644 index 00000000..3c9dab9d --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/swpc/collector.go @@ -0,0 +1,211 @@ +// ©AngelaMos | 2026 +// collector.go + +package swpc + +import ( + "context" + "encoding/json" + "log/slog" + "time" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +const ( + Name = "swpc" + defaultFastInterval = time.Minute + defaultSlowInterval = 3 * time.Hour + + keyPlasma = "swpc:plasma" + keyMag = "swpc:mag" + keyKp = "swpc:kp" + keyXray = "swpc:xray" + keyAlerts = "swpc:alerts" +) + +type Fetcher interface { + FetchPlasma(ctx context.Context) ([]PlasmaTick, error) + FetchMag(ctx context.Context) ([]MagTick, error) + FetchKp(ctx context.Context) ([]KpTick, error) + FetchXray(ctx context.Context) ([]XrayTick, error) + FetchAlerts(ctx context.Context) ([]AlertItem, error) +} + +type Ring interface { + Push(ctx context.Context, key string, score int64, payload []byte) 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 { + FastInterval time.Duration + SlowInterval time.Duration + Fetcher Fetcher + Ring Ring + Emitter Emitter + State StateRecorder + Logger *slog.Logger +} + +type Collector struct { + cfg CollectorConfig + logger *slog.Logger +} + +func NewCollector(cfg CollectorConfig) *Collector { + if cfg.FastInterval <= 0 { + cfg.FastInterval = defaultFastInterval + } + if cfg.SlowInterval <= 0 { + cfg.SlowInterval = defaultSlowInterval + } + 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 { + fast := time.NewTicker(c.cfg.FastInterval) + defer fast.Stop() + slow := time.NewTicker(c.cfg.SlowInterval) + defer slow.Stop() + + c.tickFast(ctx) + c.tickSlow(ctx) + + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-fast.C: + c.tickFast(ctx) + case <-slow.C: + c.tickSlow(ctx) + } + } +} + +func (c *Collector) tickFast(ctx context.Context) { + pushed := int64(0) + hadError := false + + if n, err := c.pushPlasma(ctx); err != nil { + c.logger.Warn("swpc plasma", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + hadError = true + } else { + pushed += n + } + if n, err := c.pushMag(ctx); err != nil { + c.logger.Warn("swpc mag", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + hadError = true + } else { + pushed += n + } + if n, err := c.pushXray(ctx); err != nil { + c.logger.Warn("swpc xray", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + hadError = true + } else { + pushed += n + } + if n, err := c.pushAlerts(ctx); err != nil { + c.logger.Warn("swpc alerts", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + hadError = true + } else { + pushed += n + } + + if pushed > 0 { + body, _ := json.Marshal(map[string]any{ + "ts": time.Now().UTC(), + "pushed": pushed, + }) + c.cfg.Emitter.Emit(events.Event{ + Topic: events.TopicSpaceWeather, + Timestamp: time.Now().UTC(), + Source: Name, + Payload: json.RawMessage(body), + }) + } + + if !hadError { + _ = c.cfg.State.RecordSuccess(ctx, Name, pushed) + } +} + +func (c *Collector) tickSlow(ctx context.Context) { + if _, err := c.pushKp(ctx); err != nil { + c.logger.Warn("swpc kp", "err", err) + _ = c.cfg.State.RecordError(ctx, Name, err.Error()) + } +} + +func (c *Collector) pushPlasma(ctx context.Context) (int64, error) { + rows, err := c.cfg.Fetcher.FetchPlasma(ctx) + if err != nil { + return 0, err + } + return pushAll(ctx, c.cfg.Ring, keyPlasma, rows, func(r PlasmaTick) int64 { return r.TimeTag.UnixMilli() }) +} + +func (c *Collector) pushMag(ctx context.Context) (int64, error) { + rows, err := c.cfg.Fetcher.FetchMag(ctx) + if err != nil { + return 0, err + } + return pushAll(ctx, c.cfg.Ring, keyMag, rows, func(r MagTick) int64 { return r.TimeTag.UnixMilli() }) +} + +func (c *Collector) pushKp(ctx context.Context) (int64, error) { + rows, err := c.cfg.Fetcher.FetchKp(ctx) + if err != nil { + return 0, err + } + return pushAll(ctx, c.cfg.Ring, keyKp, rows, func(r KpTick) int64 { return r.TimeTag.UnixMilli() }) +} + +func (c *Collector) pushXray(ctx context.Context) (int64, error) { + rows, err := c.cfg.Fetcher.FetchXray(ctx) + if err != nil { + return 0, err + } + return pushAll(ctx, c.cfg.Ring, keyXray, rows, func(r XrayTick) int64 { return r.TimeTag.UnixMilli() }) +} + +func (c *Collector) pushAlerts(ctx context.Context) (int64, error) { + rows, err := c.cfg.Fetcher.FetchAlerts(ctx) + if err != nil { + return 0, err + } + return pushAll(ctx, c.cfg.Ring, keyAlerts, rows, func(r AlertItem) int64 { return r.IssueDatetime.UnixMilli() }) +} + +func pushAll[T any](ctx context.Context, ring Ring, key string, rows []T, score func(T) int64) (int64, error) { + pushed := int64(0) + for _, r := range rows { + s := score(r) + if s == 0 { + continue + } + body, _ := json.Marshal(r) + if err := ring.Push(ctx, key, s, body); err != nil { + return pushed, err + } + pushed++ + } + return pushed, nil +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/swpc/collector_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/swpc/collector_test.go new file mode 100644 index 00000000..11bedc5f --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/swpc/collector_test.go @@ -0,0 +1,159 @@ +// ©AngelaMos | 2026 +// collector_test.go + +package swpc_test + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/swpc" + "github.com/carterperez-dev/monitor-the-situation/backend/internal/events" +) + +type fakeFetcher struct { + plasma []swpc.PlasmaTick + mag []swpc.MagTick + kp []swpc.KpTick + xray []swpc.XrayTick + alerts []swpc.AlertItem + err error +} + +func (f *fakeFetcher) FetchPlasma(_ context.Context) ([]swpc.PlasmaTick, error) { + return f.plasma, f.err +} +func (f *fakeFetcher) FetchMag(_ context.Context) ([]swpc.MagTick, error) { return f.mag, f.err } +func (f *fakeFetcher) FetchKp(_ context.Context) ([]swpc.KpTick, error) { return f.kp, f.err } +func (f *fakeFetcher) FetchXray(_ context.Context) ([]swpc.XrayTick, error) { + return f.xray, f.err +} +func (f *fakeFetcher) FetchAlerts(_ context.Context) ([]swpc.AlertItem, error) { + return f.alerts, f.err +} + +type fakeRing struct { + mu sync.Mutex + pushes map[string]int +} + +func (r *fakeRing) Push(_ context.Context, key string, _ int64, _ []byte) error { + r.mu.Lock() + defer r.mu.Unlock() + if r.pushes == nil { + r.pushes = make(map[string]int) + } + r.pushes[key]++ + return nil +} + +func (r *fakeRing) PushCount(key string) int { + r.mu.Lock() + defer r.mu.Unlock() + return r.pushes[key] +} + +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_FastTickPushesToRingsAndEmits(t *testing.T) { + now := time.Now().UTC() + ftch := &fakeFetcher{ + plasma: []swpc.PlasmaTick{{TimeTag: now, Density: "2.94", Speed: "450", Temperature: "93030"}}, + mag: []swpc.MagTick{{TimeTag: now, Bt: "5.6"}}, + xray: []swpc.XrayTick{{TimeTag: now, Flux: 1e-7, Energy: "0.1-0.8nm"}}, + alerts: []swpc.AlertItem{{ProductID: "TIIA", IssueDatetime: now, Message: "test alert"}}, + kp: []swpc.KpTick{{TimeTag: now, Kp: 3.0}}, + } + ring := &fakeRing{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := swpc.NewCollector(swpc.CollectorConfig{ + FastInterval: 20 * time.Millisecond, + SlowInterval: 50 * time.Millisecond, + Fetcher: ftch, + Ring: ring, + Emitter: emt, + State: st, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 80*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + require.GreaterOrEqual(t, ring.PushCount("swpc:plasma"), 1) + require.GreaterOrEqual(t, ring.PushCount("swpc:mag"), 1) + require.GreaterOrEqual(t, ring.PushCount("swpc:xray"), 1) + require.GreaterOrEqual(t, ring.PushCount("swpc:alerts"), 1) + require.GreaterOrEqual(t, ring.PushCount("swpc:kp"), 1) + + require.GreaterOrEqual(t, emt.Count(), 1) + for _, ev := range emt.events { + require.Equal(t, events.TopicSpaceWeather, ev.Topic) + } + require.Greater(t, st.successes, 0) + require.Equal(t, 0, st.failures) +} + +func TestCollector_FetchErrorsRecordsState(t *testing.T) { + ftch := &fakeFetcher{err: errors.New("upstream 503")} + ring := &fakeRing{} + emt := &fakeEmitter{} + st := &recordingState{} + + c := swpc.NewCollector(swpc.CollectorConfig{ + FastInterval: 20 * time.Millisecond, + SlowInterval: 50 * time.Millisecond, + Fetcher: ftch, + Ring: ring, + Emitter: emt, + State: st, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Millisecond) + defer cancel() + _ = c.Run(ctx) + + require.Equal(t, 0, ring.PushCount("swpc:plasma")) + require.Greater(t, st.failures, 0) +}