feat(monitor/collectors/swpc): collector with 1m + 3h cadences and ring-buffered persistence
This commit is contained in:
parent
608a7937e7
commit
05ca2f529f
|
|
@ -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
|
||||
}
|
||||
|
|
@ -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)
|
||||
}
|
||||
Loading…
Reference in New Issue